Skip to main content
这一页把 Task 与 CallCall Builder 与 Writer Plugin 落成代码。前四个示例用 blockx-py(blockx 仓库 python/pyproject.toml 锁定 v0.1.72),最后一节给出不依赖 SDK、直接调 WorkerService 的 Go 写法。

前置

快速开始 启动 worker 和 examples/local_test_service/local_proxy.py。示例 1 在本地就能跑通;示例 2 到 4 要读写 BlockDB 表,需要把 BLOCKDB_*_ADDR 指向可访问的 BlockDB,并且不设 STUB_BLOCKDB_DELAY_MS
blockx-py 只通过 PROXY_SOCKET_PATH 指定的 UDS proxy 连接。没设这个变量会直接 RuntimeError,不会静默连到别处。

示例 1:一次性计算并取回结果

InputsCallConfig + ReturnValueHandler。这就是 examples/local_test_service/submit.py
要点:
  • LocalTestService 是 worker 内置的示例能力后端,本地不需要 BlockDB。换成任何自包含的纯计算函数也一样能跑。
  • 结果按各次调用的完成序收集,不保证与 callList 行序一致。需要对应关系时让返回值带上入参。
  • ReturnValueHandler 的 task 被 SDK 路由到 block 集群,因为 bundle worker 没有注册这个 handler。

示例 2:按区块计算并写入 block 表

BlockTableCallConfig + BlockTableWriteHandler(block=...)。订阅源表的新区块,每个区块提交一个 task:
要点:
  • 传 callable 时 SDK 用 inspect.getsource 抓函数体并补一行 _ = trace_to_token_transfer;也可以直接传 Function(source_code=...) 或注册函数 ID 字符串。
  • params 决定函数的位置参数:这里 chain_id="eth"record 是源表的一行 dict。
  • BlockTableWriteHandler 只投影返回 dict 中属于目标表的列,written_rows 是这次写入的行数。
  • triggerSources 里可以加 "operator": filter(...) 在 Builder 阶段先过滤、去重,减少 call 数量。

示例 3:写入普通表

把示例 2 的 handler 换成 NormalTableWriteHandler,目标就从 block 表变成普通表(L1)。它不需要 block
condition.column 不能是主键 id,非法值在构造时就抛 ValueError。这个 handler 不带区块上下文,同一个实例可以直接交给示例 4 的回填。

示例 4:按 bundle 回填

历史回填按 bundle(1000 个区块一段)提交 task,用 BlockBundleCallConfig 读 bundle 分区的 parquet,用 BlockTableWriteHandler(block_bundle=...) 覆盖写目标表的同一个 bundle。backfill_block_bundles 帮你按区间逐 bundle 实例化和提交:
只发一个 bundle 时直接构造带号码的 config 和 handler:
BlockTableWriteHandler 的形状必须跟集群走:bundle 集群要 block_bundle=,block 集群要 block=。发错形状服务端要等所有 call 算完才在 Writer 阶段拒绝,所以 SDK 在 submit() 前就校验并抛 ValueError
日常使用推荐 blockx-py 的上层封装 Pipeline:一份 triggers + target_tablebackfill(block_start, block_end) 做历史回填,update() 先补齐缺口再转实时订阅。bundle 换算、集群路由、slot 退避和多表对齐都由它处理:

读结果

task.submit() 返回 TaskResult retryable 决定是否重投,不要从错误文本推断:
Task.task_id 每次 build 都是新的 uuid;排查时把它记到日志里,worker 的 task finished 日志按它检索。

不用 SDK:直接调 gRPC

Go 侧的贡献者经常需要在测试或工具里直接对 Worker 提交 task。协议是 api/grpc/worker/worker.protoWorkerService,Go 绑定在 api/grpc/worker/workerpb。下面的程序对本地 worker 走 RequestTaskSlot → SubmitTask → WatchTasks 三步,提交的是示例 1 同样的 task:
  • RequestTaskSlot 可以省略:SubmitTask 不带 slot_id 时 Worker 会内联申请一次,没有容量返回 ResourceExhausted
  • 经 Coordinator 时先调 coordinator.v1.CoordinatorService.ReserveWorkerSlot(task_id)api/grpc/coordinator/coordinator.proto),再用返回的 worker_addr 建连接、用 slot_id 提交。
  • 同一 task_id 重复 SubmitTask 是幂等的:task 仍在跑返回 RUNNING,结果仍在保留窗口内返回终态。
  • 测试里可以直接用 internal/testutilRequestTaskSlot / SubmitTask / PollUntilTerminal 辅助函数,它们封装的就是上面的调用。

常见错误

相关文档

  • 快速开始:本地起 worker 与 proxy 的最短路径。
  • Task 与 CallTaskInput 字段、状态与失败码。
  • Call Builder 与 Writer Plugin:每种 config 的 wire 形状。
  • 协议与接口WorkerService / CoordinatorService 的完整字段。
  • 测试:process / system E2E 里的提交辅助函数。
  • blockx-py 仓库 docs/api_reference.md:SDK 全部参数与 Pipeline 的行为细节。