Skip to main content
task 通过两个 (type, config) 对声明前后两端:functionCallConfig 选一个 Call Builder,resultHandler 选一个 Writer Plugin(blockx-py 里叫 Handler)。type 必须是当前 Worker 进程注册过的名字;config 由对应实现自己解析,Worker 在提交时只校验它是合法 JSON。

Call Builder

BlockTableCallConfig

  • 一个 task 可以有多个 triggerSources,各自独立扫描后合并成一份 call 列表。
  • functionIdcode.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(单数),与 BlockTableCallConfigfunctionId / params 不同。
  • bundle worker 默认开启 stream build:扫描与执行流水线化,百万行也不必先物化整份 call 列表。

Operator:过滤与去重

operator 是一个 Substrait 风格的 relations 链:恰好一个 read,后接任意个 filterdeduplicate,用 input.relation_id 串成线性链。
  • filter 支持 andornotequalnot_equalgtgteltlteis_nullis_not_null,按 SQL 三值逻辑求值,只保留结果为 TRUE 的行。
  • deduplicate 按指定列去重,key 用带类型标签的确定性编码。
  • dbScan 在 Worker 内存里执行这条链;bundleScan 把它翻译成 SQL 下推给 DuckDB。
  • blockdb-py 的 filter() / Operator 直接产出这个结构,可以原样放进 operator
格式细节见 blockx 仓库 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_FAILEDretryable 由插件的错误分类决定。

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

  • 按主键 id upsert。
  • condition 可选:{"column": "ts", "policy": "UpdateIfSmaller"} 表示只有新值在该列上更小才覆盖;policyUpdateIfSmallerUpdateIfLarger。乱序或重放写入下天然幂等。
  • 不带区块上下文,同一份配置可以在实时和回填两条路径里复用。

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 回传。
  • BlockTableWriteHandlerconfig 形状必须跟集群走:block 集群带 block,bundle 集群带 blockBundle
  • operator 在 Builder 阶段过滤、去重触发行,减少 call 数量。
继续阅读: