blockx 负责算——把「源表 + 函数 + 目标表」打包成任务提交给远端 worker。日常写索引管道用 Pipeline 就够了;这一页是 Pipeline 内部实际在用的底层 API,给需要绕开 Pipeline、做精细控制的场景。
心智模型
一个任务 = 一个 CallConfig(算什么)+ 一个 ResultHandler(结果怎么处理)。Pipeline 就是把这一整套包起来的高层封装。
Trigger 与 func 的三种给法
params 里除了 SOURCE_ROW,其余原样作为函数的位置参数——只有位置参数,没有命名参数,这是跨语言设计决定的。
CallConfig(现有三种)
ResultHandler(现有四种)
组装 + 提交
success 只代表 worker 受理了这个任务,写库是后台异步做的——这一刻数据不一定已经进表。真实进度看目标表的 _write 子表。
task.submit() 常用参数:
Function 互调用
一个 Function 可以在执行过程中调用另一个已保存的 Function:
在 worker 里跑走同一个 task 内的子调用,在本地 / Notebook 里跑走 sync-invoker。
Operator:执行前预过滤 / 去重
有些过滤、去重逻辑不必塞进业务函数——可以在 Trigger 上挂一条 operator 管道,引擎会在把行喂给函数之前先做掉:
operator 接受 Operator 对象、dict,或者它们的 JSON 字符串,三种写法效果一样。
调试:不落库预演
联调新转换逻辑时,用环境变量开启 debug 模式——worker 端不会真正写目标表,读 / 算逻辑照常跑,只有最后一步落库被替换成打日志:
这个开关和 blockdb-py 共用,设一次就把整条 SDK 链路一起切到 debug,代码不用改。上线前记得关掉,否则生产任务的写入会全部进日志而不落库。