(type, config) 对声明前后两端:functionCallConfig 选一个 Call Builder,resultHandler 选一个 Writer Plugin(blockx-py 里叫 Handler)。type 必须是当前 Worker 进程注册过的名字;config 由对应实现自己解析,Worker 在提交时只校验它是合法 JSON。
Call Builder
BlockTableCallConfig
- 一个 task 可以有多个
triggerSources,各自独立扫描后合并成一份 call 列表。 functionId与code.sourceCode二选一:前者是注册函数,后者是 inline 源码。params为空时整行作为唯一参数;params非空但没有"${表名}"占位符时,所有 call 共享同一份静态参数,因此只会执行一次。operator可选,省略时触发表在该区块上的每一行都生成一个 call。
InputsCallConfig
- 不做任何 IO。
callList[i]直接成为第 i 个 call 的参数数组。 - 字段名是
function,不是functionId;blockx-py 会自动做这层重命名。 - 结果按完成序汇聚,不保证与
callList的行序一致。需要对应关系时,让返回值带上入参。
BlockBundleCallConfig
- bundle 是按 1000 个区块为一段切分的数据分区,由
blockBundle.number标识。 tableId只报源表身份。Worker 按(tableId, number)从该表的 Iceberg 快照解析这个 bundle 的数据文件;state 表经<tableId>._archive解析。- 字段名是
function/param(单数),与BlockTableCallConfig的functionId/params不同。 - bundle worker 默认开启 stream build:扫描与执行流水线化,百万行也不必先物化整份 call 列表。
Operator:过滤与去重
operator 是一个 Substrait 风格的 relations 链:恰好一个 read,后接任意个 filter 和 deduplicate,用 input.relation_id 串成线性链。
filter支持and、or、not、equal、not_equal、gt、gte、lt、lte、is_null、is_not_null,按 SQL 三值逻辑求值,只保留结果为TRUE的行。deduplicate按指定列去重,key 用带类型标签的确定性编码。- dbScan 在 Worker 内存里执行这条链;bundleScan 把它翻译成 SQL 下推给 DuckDB。
- blockdb-py 的
filter()/Operator直接产出这个结构,可以原样放进operator。
docs/specs/substrait_ast.md。
Writer Plugin
共同行为:
- 写表插件先读目标表的列定义,只投影返回 dict 中属于目标表的列;
outputs里的null被跳过。 - 任意 handler 的
config里带"debug": true时,写表插件只统计行数并打table_write_debug日志,不读 schema、不发写请求。blockx-py 用环境变量BLOCKDB_DEBUG=1打开它。 - 不带
resultHandler也是合法的:Calls 阶段成功后 task 直接进入终态,outputs被丢弃。 - Writer Plugin 是 fail-fast:写入失败即 task
PLUGIN_FAILED,retryable由插件的错误分类决定。
BlockTableWriteHandler
- 一个名字对应两套实现,由 Worker 进程的部署形态(
block/bundle)在启动时装配,调用方不能选。 - block 集群的实现要求
config.block,一次 task 一次写入,返回{written_rows, finish_time}。 - bundle 集群的实现要求
config.blockBundle,经 BlockDB BatchWrite job(Init → presigned PUT → Commit)覆盖写目标表的同一个 bundle。 - 形状不匹配时在 Writer 阶段以
PLUGIN_FAILED失败(block deployment requires config.block/bundle deployment requires config.blockBundle),此时 call 已经全部执行过。blockx-py 在提交前先校验形状与目标集群,避免白算一遍。
NormalTableWriteHandler
- 按主键
idupsert。 condition可选:{"column": "ts", "policy": "UpdateIfSmaller"}表示只有新值在该列上更小才覆盖;policy取UpdateIfSmaller或UpdateIfLarger。乱序或重放写入下天然幂等。- 不带区块上下文,同一份配置可以在实时和回填两条路径里复用。
ReturnValueHandler
- 把
outputs原样拼成一个 JSON 数组放进pluginResults[i].result。blockx-py 里读result.handler_result。 - bundle worker 没有注册它,带它的 task 提交到 bundle 集群会在
SubmitTask时被InvalidArgument拒绝;blockx-py 因此把这类 task 一律路由到 block 集群。
组合速查
小结
- Call Builder 决定输入:
BlockTableCallConfig读区块、InputsCallConfig用自带参数、BlockBundleCallConfig读 bundle。 - Writer Plugin 决定结果去向:
BlockTableWriteHandler写 L2、NormalTableWriteHandler写 L1、ReturnValueHandler回传。 BlockTableWriteHandler的config形状必须跟集群走:block 集群带block,bundle 集群带blockBundle。operator在 Builder 阶段过滤、去重触发行,减少 call 数量。
- Plugin 系统:接口、注册表、装配与扩展一个新的 Builder / Plugin。
- Bundle 集群:bundleScan 流式扫描与 BatchWrite 写路径。
- 提交 task 示例:每一种组合的 blockx-py 写法。