跳轉到主要內容
ChainStream 透過 Kafka Streams 提供多鏈實時鏈上資料流。相較於 GraphQL Subscriptions 和 WebSocket,Kafka Streams 面向對延遲敏感、可靠性要求高的服務端應用場景,提供更低延遲、更強容錯的資料消費能力。

Protobuf Schema 倉庫

ChainStream 官方 Protobuf Schema 定義,支援 Go 和 Python,包含 EVM、Solana、TRON 所有訊息型別。

支援矩陣

所有鏈還支援 token-suppliestoken-pricestoken-holdingstoken-market-capstrade-stats 等 Topics。詳見完整 Topic 列表。

Kafka Streams vs WebSocket 選型指南

何時選擇 Kafka Streams

延遲敏感

延遲是首要考量,應用部署在雲端或專用伺服器

訊息可靠

不可接受丟失任何訊息,需要持久可靠的資料消費

複雜處理

需要對資料做複雜計算、過濾或格式化,超出預處理能力範圍

水平擴充套件

需要多例項水平擴充套件消費能力

何時選擇 WebSocket

快速原型

正在構建原型,開發速度是首要因素

統一介面

應用同時需要歷史資料和實時資料,需要統一查詢與訂閱介面

瀏覽器端

應用直接在瀏覽器端消費資料(Kafka Streams 僅支援服務端)

動態過濾

需要根據頁面內容動態過濾資料

對比總結


接入憑證獲取

Kafka Streams 使用獨立的認證憑證,需要聯絡 ChainStream 團隊申請開通。
1

聯絡申請

傳送郵件至 support@chainstream.io 申請 Kafka Streams 接入許可權
2

獲取憑證

稽核透過後,您將收到以下憑證資訊:
  • Username
  • Password
  • Broker 地址列表
3

配置連線

使用獲取的憑證配置 Kafka 客戶端連線

連線配置

Broker 地址

Broker 地址將在您的申請稽核透過後,隨憑證資訊一同提供。請勿使用任何未經授權的地址進行連線。

SASL_SSL 連線配置


Topic 命名規範與完整列表

命名規範

Topic 命名遵循以下 pattern:
其中 {chain} 包括:solbscethtron

訊息型別說明

完整 Topic 列表

以下 Topics 適用於所有支援的鏈(將 {chain} 替換為 solbsceth):
完整的 Protobuf Schema 和 Topic 對映請參考 streaming_protobuf 倉庫

消費模式與 Offset 管理

訂閱 topic 時需要關注兩個核心配置:

Offset 策略選擇

消費者在連線 Kafka 後,需要決定從哪個位置開始讀取訊息。兩種常見策略:
每次連線從當前最新位置開始,適合只關心實時資料的場景。重連後不會回溯歷史訊息。

Group ID 規則

多例項部署同一 Group ID 可實現故障轉移和負載均衡——同一 topic 的訊息只會被 Group 中的一個例項消費,Kafka 自動在例項間分配分割槽。
通常建議每個 topic 對應一個獨立 consumer,因為不同 topic 的訊息解析邏輯不同。

Quick Start:5 分鐘跑通第一個 Consumer

以下示例展示如何消費 eth.dex.trades topic 並解析 DEX 交易資料。
1

獲取 Protobuf Schema

從官方倉庫克隆 Schema 定義:
或作為 Git submodule 新增到專案:
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 時需要注意以下訊息特性:
Stream 不做預過濾,包含 topic 內的所有訊息和完整資料。這意味著消費端需要有足夠的網路吞吐、伺服器效能和高效的解析程式碼。
同一個代幣或同一個賬號的訊息嚴格按 block 順序到達。這意味著針對特定代幣或錢包地址的事件流是有序的,方便追蹤狀態變化。但不同代幣/賬號之間的訊息到達順序不做保證。
同一條訊息可能被投遞多次。如果重複處理會造成問題,消費端需要維護快取或儲存來實現冪等處理。
ChainStream 保證每條訊息的完整性,訊息不會被拆分。無論區塊包含多少交易,您收到的訊息都是完整的資料單元。
訊息使用 Protobuf 編碼,比 JSON 更緊湊。消費端需要使用對應語言的 Protobuf 庫進行解析。

延遲模型

Kafka Streams 的延遲取決於資料在管道中經過的處理環節。同一條鏈的不同 topic 延遲不同:

Broadcasted vs Committed

處理管道延遲

資料從區塊鏈節點到 Kafka topic 的每一層轉換(解析、結構化、enrichment)都會引入約 100-1000ms 的延遲:
  • raw topic:延遲最低,接近原始節點資料
  • transactions topic:經過解析和結構化
  • dextrades topic:延遲相對更高,但資料更豐富
如果延遲是首要考量,優先選擇離原始資料最近、你能有效解析的 topic。

最佳實踐

分割槽並行消費

Kafka topic 被劃分為多個分割槽(partition),每個分割槽需要並行讀取以最大化吞吐量。 訊息的分割槽鍵設定為 代幣地址錢包地址(所有鏈統一),這確保:
  • 同一代幣的所有事件路由到同一分割槽,保證順序性
  • 同一錢包的所有餘額變動路由到同一分割槽,方便狀態追蹤
建議為每個分割槽分配一個獨立執行緒,確保負載均衡。

持續消費,不阻塞主迴圈

Consumer 的讀取迴圈應保持持續執行,避免因訊息處理阻塞而導致積壓。如果需要對訊息做處理,應採用非同步處理模式:主迴圈負責讀取,處理邏輯委託給 worker 執行緒。

訊息處理效率

批次處理可以降低開銷,但需要在批次大小和延遲之間權衡。在 Go 中可以使用 channel + worker group 實現併發處理。

鏈特定文件

EVM Streams

Ethereum、BSC、Base、Polygon、Optimism

Solana Streams

Solana 高吞吐資料流

TRON Streams

TRON 網路資料流

相關文件

實時資料流

WebSocket 實時資料接入指南

認證指南

獲取 Access Token