メインコンテンツへスキップ
ChainStream は Kafka Streams を通じてマルチチェーンのリアルタイムオンチェーンデータストリームを提供します。GraphQL Subscriptions や WebSocket と比較して、Kafka Streams はレイテンシに敏感で高信頼性が求められるサーバーサイドアプリケーションシナリオ向けに設計されており、より低いレイテンシと高い耐障害性を備えたデータ消費を提供します。

Protobuf スキーマリポジトリ

公式 ChainStream Protobuf スキーマ定義。Go と Python をサポートし、EVM、Solana、TRON のすべてのメッセージタイプを含みます。

サポートマトリクス

すべてのチェーンで token-suppliestoken-pricestoken-holdingstoken-market-capstrade-stats Topic もサポートしています。詳細は完全な Topic リストを参照してください。

Kafka Streams vs WebSocket 選択ガイド

Kafka Streams を選ぶべきとき

レイテンシ重視

レイテンシが最重要課題で、アプリケーションがクラウドまたは専用サーバーにデプロイされている場合

メッセージの信頼性

メッセージの欠落が許容できず、耐久性があり信頼性の高いデータ消費が必要な場合

複雑な処理

前処理機能を超える複雑な計算、フィルタリング、フォーマットが必要な場合

水平スケーリング

消費能力のためにマルチインスタンスの水平スケーリングが必要な場合

WebSocket を選ぶべきとき

高速プロトタイピング

プロトタイプの構築で、開発速度が最優先の場合

統一インターフェース

アプリケーションが履歴データとリアルタイムデータの両方を統一されたクエリ・サブスクリプションインターフェースで必要とする場合

ブラウザサイド

アプリケーションがブラウザで直接データを消費する場合(Kafka Streams はサーバーサイドのみサポート)

動的フィルタリング

ページコンテンツに基づいてデータを動的にフィルタリングする必要がある場合

比較まとめ


認証情報の取得

Kafka Streams は独立した認証情報を使用し、ChainStream チームに連絡してアクセスを申請する必要があります。
1

申請連絡

support@chainstream.io にメールを送信して Kafka Streams へのアクセスを申請
2

認証情報の受領

承認後、以下の認証情報を受け取ります:
  • ユーザー名
  • パスワード
  • ブローカーアドレスリスト
3

接続の設定

受け取った認証情報を使用して Kafka クライアントの接続を設定

接続設定

ブローカーアドレス

ブローカーアドレスは申請承認後に認証情報と一緒に提供されます。未許可のアドレスでの接続は行わないでください。

SASL_SSL 接続設定


Topic 命名規則と完全リスト

命名規則

Topic は以下の命名パターンに従います:
{chain} には solbscethtron が含まれます。

メッセージタイプ

完全な Topic リスト

以下の Topic はすべてのサポートチェーンに適用されます({chain}solbsceth に置き換え):
完全な Protobuf スキーマと Topic マッピングについては、streaming_protobuf リポジトリを参照してください。

消費モードとオフセット管理

Topic をサブスクライブする際に考慮すべき 2 つのコア設定:

オフセット戦略の選択

コンシューマーは Kafka に接続後、どこからメッセージの読み取りを開始するかを決定する必要があります。2 つの一般的な戦略:
各接続時に現在の最新位置から開始。リアルタイムデータのみを気にするシナリオに適しています。再接続時に過去のメッセージのリプレイはありません。

Group ID ルール

同じ Group ID で複数インスタンスをデプロイすると、フェイルオーバーと負荷分散が可能になります — 同じ Topic のメッセージは Group 内の 1 つのインスタンスのみが消費し、Kafka がインスタンス間でパーティションを自動分配します。
異なる Topic は異なるメッセージパースロジックを持つため、各 Topic に独立したコンシューマーを持つことを推奨します。

クイックスタート:5 分で最初のコンシューマー

以下の例は、eth.dex.trades Topic を消費して DEX 取引データをパースする方法を示しています。
1

Protobuf スキーマの取得

公式リポジトリからスキーマ定義をクローン:
またはプロジェクトに Git サブモジュールとして追加:
2

依存関係のインストール

3

設定と消費


コアデータ構造

すべてのメッセージタイプは以下のベース構造を共有しています(common/common.proto で定義):

ベース構造

ブロック情報:

主要メッセージタイプ

Topic: {chain}.dex.trades
Trade コアフィールド:TradeProcessed 拡張フィールド(processed topic):
Topic: {chain}.tokens{chain}.tokens.created
Token コアフィールド:
Topic: {chain}.balances
Balance コアフィールド:
Topic: {chain}.dex.pools
DexPool コアフィールド:
Topic: {chain}.candlesticks
Topic: {chain}.trade-stats
Topic: {chain}.token-holdings
Topic: {chain}.token-prices
Topic: {chain}.token-supplies
完全な Protobuf 定義については streaming_protobuf リポジトリを参照してください。

メッセージ特性と注意事項

Kafka Streams を消費する際、開発者は以下のメッセージ特性を認識しておく必要があります:
ストリームは事前フィルタリングを行わず、Topic 内のすべてのメッセージと完全なデータを含みます。コンシューマーには十分なネットワークスループット、サーバー性能、効率的なパースコードが必要です。
同じトークンまたは同じアカウントのメッセージは厳密にブロック順序で到着します。特定のトークンやウォレットアドレスのイベントストリームは順序が保証され、状態変更の追跡が容易です。ただし、異なるトークン/アカウント間のメッセージ到着順序は保証されません。
同じメッセージが複数回配信される場合があります。重複処理が問題を引き起こす場合、コンシューマーはキャッシュやストレージを使用した冪等処理を実装する必要があります。
ChainStream は各メッセージの完全性を保証します。メッセージが分割されることはありません。ブロックに含まれるトランザクション数に関係なく、受信するメッセージは完全なデータ単位です。
メッセージは Protobuf エンコーディングを使用し、JSON よりコンパクトです。コンシューマーは対応する言語の Protobuf ライブラリを使用してパースする必要があります。

レイテンシモデル

Kafka Streams のレイテンシは、データがパイプライン内で通過する処理ステージに依存します。同じチェーンの異なる Topic は異なるレイテンシを持ちます:

Broadcasted vs Committed

パイプラインレイテンシ

ブロックチェーンノードから Kafka Topic までの各変換レイヤー(パーシング、構造化、エンリッチメント)がおよそ 100〜1000ms のレイテンシを追加します:
  • raw topic: 最低レイテンシ、生のノードデータに最も近い
  • transactions topic: パース・構造化済み
  • dextrades topic: 相対的に高いレイテンシだが、よりリッチなデータ
レイテンシが最優先の場合は、効果的にパースできる生データに最も近い Topic を選択してください。

ベストプラクティス

パーティション並列消費

Kafka Topic は複数のパーティションに分割されており、各パーティションの並列読み取りがスループットを最大化します。 メッセージパーティションキーはトークンアドレスまたはウォレットアドレス(全チェーンで統一)に設定されており、以下が保証されます:
  • 同じトークンのすべてのイベントが同じパーティションにルーティングされ、順序を保証
  • 同じウォレットのすべての残高変更が同じパーティションにルーティングされ、状態追跡が容易
各パーティションに独立したスレッドを割り当てて負荷分散することを推奨します。

継続的消費、メインループをブロックしない

コンシューマーの読み取りループは継続的に実行し、メッセージ処理のブロックによるバックログを避けてください。メッセージの処理が必要な場合は非同期処理モードを採用:メインループは読み取りを担当し、処理ロジックはワーカースレッドに委任します。

メッセージ処理効率

バッチ処理はオーバーヘッドを削減できますが、バッチサイズとレイテンシのバランスが必要です。Go では channel + ワーカーグループを使用した並行処理が効果的です。

チェーン別ドキュメント

EVM Streams

Ethereum、BSC、Base、Polygon、Optimism

Solana Streams

Solana 高スループットデータストリーム

TRON Streams

TRON ネットワークデータストリーム

関連ドキュメント

リアルタイムストリーミング

WebSocket リアルタイムデータ統合ガイド

認証ガイド

Access Token の取得