Skip to main content
pipeline 是一个 Go 库,不是独立进程。它被编译进执行客户端,通过 EVM 追踪钩子采集区块执行数据,序列化后写入 S3,并向 Kafka 发布区块变更通知。

模块

两种集成模式

追踪器嵌入区块执行流程,实时采集并上传。零额外延迟,但需要修改执行客户端的核心代码。

采集机制

CallTracer

维护调用栈并构建调用树。OnEnter 压栈,OnExit 出栈并挂到父帧下。
  • 存储变更标记:识别 SSTORE 指令,把 StorageChange 向上传播到祖先调用
  • 失败隔离:调用失败时标记自身及所有子调用的 ParentFailed,展平后进入 ErrorTraces / ErrorEvents
  • ID 生成:trace ID 为 hash(tx_id, parent_trace_id, position),event ID 为 hash(parent_trace_id, position),保证跨区块可复现

StateDiff 的两条路径

第二种用于没有 commit hook 的客户端。originRoot == root 时不生成差异对象。

上传与发布

OnCommit 里并发执行四路上传:
StateDiff 用 RLP 是因为它只被内部组件消费,体积比 JSON 小得多;BlockFile 用 JSON 是因为外部消费者需要能直接解析。 配置 s3_temp_dir 后启用本地缓存:先落盘再异步上传,上传成功后删除本地文件。S3 返回 5xx 时按退避重试,计入 pipeline/s3_upload_retry

Leader 选举

多个写节点可以同时运行 pipeline,但只有 Leader 向 Kafka 发布通知,所有实例都上传 S3。S3 对象以哈希和 state root 为键,重复上传是幂等的。
配置 etcd_endpoints,把 is_backup 留空。
键不存在时抢占,存在时 watch。键被删除后随机退避再抢,避免惊群。成为 Leader 前等待 grace_period(默认 10 秒),让上一任 Leader 完成收尾。
配置了 etcd 时还会启用写节点注册:以租约方式写入 {chainID}[/{version}]/writers/{nodeID},节点故障时自动清理。nodex-proxy 的 /{chainId}/writers 管理接口读的就是这份数据。

配置

传给 NewPipelineTracer 的 JSON(即 geth 的 --vmtrace.jsonconfig):

指标

通过 go-ethereum 的 metrics 体系暴露。

开发

相关文档

仓库内还有两份更细的文档:

接口契约

Kafka 消息与 S3 键的精确格式。

接入新链

把 pipeline 适配到 Geth、旧版 Geth 或 Reth 分叉。