python/pyproject.toml 锁定 v0.1.72),最后一节给出不依赖 SDK、直接调 WorkerService 的 Go 写法。
前置
- 本地 worker
- 线上环境
按 快速开始 启动 worker 和
examples/local_test_service/local_proxy.py。示例 1 在本地就能跑通;示例 2 到 4 要读写 BlockDB 表,需要把 BLOCKDB_*_ADDR 指向可访问的 BlockDB,并且不设 STUB_BLOCKDB_DELAY_MS。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 实例化和提交:
BlockTableWriteHandler 的形状必须跟集群走:bundle 集群要 block_bundle=,block 集群要 block=。发错形状服务端要等所有 call 算完才在 Writer 阶段拒绝,所以 SDK 在 submit() 前就校验并抛 ValueError。Pipeline:一份 triggers + target_table,backfill(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.proto 的 WorkerService,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/testutil的RequestTaskSlot/SubmitTask/PollUntilTerminal辅助函数,它们封装的就是上面的调用。
常见错误
相关文档
- 快速开始:本地起 worker 与 proxy 的最短路径。
- Task 与 Call:
TaskInput字段、状态与失败码。 - Call Builder 与 Writer Plugin:每种 config 的 wire 形状。
- 协议与接口:
WorkerService/CoordinatorService的完整字段。 - 测试:process / system E2E 里的提交辅助函数。
- blockx-py 仓库
docs/api_reference.md:SDK 全部参数与Pipeline的行为细节。