Skip to main content
Plugin 系统是 Worker 进程内的两组可插拔实现:Call Builder 在 Executor 阶段之前把 task 配置变成 CallListWriter Plugin(task 侧叫 ResultHandler)在 Executor 阶段之后接收所有 call 的返回值并落库或透传。它们对应 架构总览 中 Task 主链路的 Builder 阶段和 Writer 阶段,由 Worker 的 orchestrator 调度。

职责与边界

  • Call Builder 只负责”读输入、产出 CallList”。它不派发 call、不拉取函数代码、不缓存结果。
  • Writer Plugin 只负责”消费 outputs、写出或返回结果”。它不参与 call 的重试或去重。
  • 一个 task 只有一个 FunctionCallConfigtype/config)和至多一个 ResultHandlertype/config)。type 是注册名,configjson.RawMessage,由对应实现自己解析。
  • Builder / Plugin 都不在构造时持有 IO;TaskIOReader / TaskIOBuild / Execute 调用时由 Worker 以参数传入,IO 限流和重试在 IO 子系统里做(见 IO 访问子系统)。
  • Worker 进程注册哪些 Builder / Plugin 由入口 app.Profile 声明。不在 Profile 里的 type 在运行期 lookup miss,task 以 BUILDER_NOT_FOUND / PLUGIN_NOT_FOUND 干净失败。
  • Writer Plugin 是 fail-fast:Execute 返回 error 即 task 失败;是否允许上游重调度由 PluginError.Retryable 决定。
  • outputs 中每个非 nil 元素必须是 json.RawMessage,内容是单个对象(单行)或对象数组(多行)。写表类插件用 types.VisitOutputRows 统一校验和遍历。
  • 流式 Builder(StreamingCallBuilder)把”静态解析”和”扫描产出”拆成两阶段:PrepareStream 失败仍是零 call 派发的原子失败;Run 中途失败则先前批次可能已执行。

代码位置

核心类型与接口

以下签名从代码复制。
event/types 里还声明了 SchemaProvider 接口,但当前没有生产调用方;实际列定义都走 schema.ReadColumns

数据流 / 执行流程

1

激活:构造 TaskCtx

Orchestrator.activateAndRunorchestrator_phases.go)从 core 取回解析好的 TaskCtx,注入 Ctx / Cancel,并把 ResultHandler.Config 里的 debug:true 提升为 TaskCtx.DebugMode
2

Builder 阶段

runBuilderPhase 先取 builderSemAdmission.BuilderSlots),按 FunctionCallConfig.TypebuilderRegistryLookup。命中 StreamingCallBuilderStreamBuild.Enabled 时走 PrepareStream,否则走 Build。成功后 releaseBuilderConfigPayloadFunctionCallConfig.Config 置 nil 释放请求 JSON。失败码:BUILDER_NOT_FOUND / BUILDER_FAILEDRetryable 取决于是否是可重试 IOError
3

Calls 阶段

非流式:runCallPhase 把整份 CallList 交给 dispatcher。流式:runCallStreamPhasescanSem 下跑 producer.RunmakeBatchEmit 提供带信用(MaxOutstandingBatches)的背压 emit。详见 Call 执行子系统
4

Writer 阶段

runWriterPhasepluginSemAdmission.PluginSlots),若 ResultHandler 非 nil 则 Lookupplugin.Execute(resultCtx, ioTaskCtx, cmd.Outputs, ioScope);返回值封装为 commontypes.PluginResult 进入 TaskResult.ExecuteResult.PluginResultsResultHandler 为 nil 时阶段直接成功。

内置 Call Builder

  • DBScanBuilderdbscan/dbscan.go):解析 DBScanBuilderDecl{Block, TriggerSources};每个 TriggerSource{Table, Operator, FunctionID, Code, Params} 一个 goroutine(errgroup.SetLimit(3))。对每个 source:校验表名标识符 → NormalizeInlineCodeOperator.Validateio.Readblockdb.ReadOp[*GetBlockRowsRequest]Op: "GetBlockRows",逻辑表名原样传给 BlockDB)→ Operator.ApplyInMemory(rows)(先 filter 后 deduplicate,按 relations 链顺序)→ 按 Params 模板拼 args。Params 中等于 "${<table>}" 的位置替换成行 map;Params 为空则整行作为唯一参数;Params 非空但没有占位符则所有 call 共享同一份静态 ArgsJSONCallID 格式为 {taskID}-{table}-{height}-{i}
  • CallListCallConfigBuildercall_list/call_list.go):解析 {function, code, callList: [][]any},不做 IO,第 i 项生成 Call{FunctionID, Args, CallID: "{taskID}-{i}"}
  • BundleScanBuilderbundlescan/):从 blockBundle.number + triggerSources[].tableId(或兼容字段 bundleBucket)解析 parquet 位置,经 DuckDB 流式读行并把 Operator 翻译成 SQL 下推。同时实现 StreamingCallBuilderstream.go,批边界 defaultBatchMaxRows = 8192 / defaultBatchMaxArgsBytes = 32 MiB),并在可证明行 key 唯一时置位 Call.NoResultCache。细节见 Bundle 集群

Operator(Plan)

TriggerSource.Operator*types.Plan,JSON 形如 {"relations": [{"read": ...}, {"filter": ...}, {"deduplicate": ...}]},relation 之间用 input.relation_id(数组下标)串成线性链。Plan.Validate 要求恰好一个 read、链无分叉、所有 selection.field 落在 read.base_schema.names 内。filter 用 SQL 三值逻辑(NULL 比较得 UNKNOWN,只保留 TRUE),equal(col, null)IS NULL 解释;deduplicate key 用带类型标签和长度前缀的确定性编码(appendDedupValue)。支持的 function_referenceandornotequalnot_equalgtgteltlteis_nullis_not_null Plan.UnmarshalJSON 还接受一种 legacy 形态:顶层带 type 字段的表达式树(Attribute / Literal / EqualTo / In / And / Or / Not 等),会被 legacyConditionToPlan 转成单个 filter relation。 完整格式见 blockx 仓库 docs/specs/substrait_ast.md

内置 Writer Plugin

所有写表插件在 TaskCtx.DebugMode 为 true 时只 PreviewOutputs 统计行数并打 table_write_debug 日志,不读 schema、不发写请求。

devstub

internal/worker/devstub/ 提供 StaticCallBuilder(名字 static,返回固定 call)、PayloadCallBuilderpayload,从 config.calls 取 call,config.builderFail 模拟失败)、LogWriterPluginlog)、TestWriterPlugintestconfig.pluginFail 模拟失败)。它们只被 cmd/worker 的 Profile 列入,cmd/bundle_worker 不注册;用于 e2e / perf(如 e2e/perf/s2_call_count_payload_test.go)和本地联调,不要在生产 task 里使用。

配置

Plugin 本身没有独立配置文件;相关项都在 internal/worker/app/config.goWorkerFullConfig

扩展点

新增一个 Call Builder

1

定义名字

internal/plugin/callbuilder/types/types.goconst 块加一个 CallBuilderName,值就是 task 里 functionCallConfig.type 的字符串。
2

实现接口

新建 internal/plugin/callbuilder/<name>/,实现 types.CallBuilder。在 Build 里用 sonic.Config{UseNumber: true} 解析 taskCtx.FunctionCallConfig.Config,用 types.NormalizeInlineCode 处理 inline code,通过传入的 io.Read 做所有读 IO。逐行 args 用 types.CompactCall 产出。若要支持流式,再实现 types.StreamingCallBuilderPrepareStream 里做完全部静态校验,Run 里按批 emit,每批换新的 backing array,遇到 types.ErrProductionStopped 原样返回)。
3

装配

internal/worker/app/app.gonewBuilderRegistryswitch name 加一个 case 构造它;再把名字加进需要它的 Profile(cmd/worker/main.gocmd/bundle_worker/main.goBuilders 列表)。不加 case 会在启动时报 unknown call builder
4

测试

包内单测(参考 dbscan/dbscan_test.gocall_list/call_list_test.go:用 fake TaskIOReader 断言请求类型、CallID、args、错误分类);internal/worker/app/profile_test.goTestNewBuilderRegistry_ProfileMembership 系列补名字;如需端到端,参考 e2e/system/logical_type/dbscan_operator_test.go

新增一个 Writer Plugin

1

定义名字

internal/plugin/event/types/types.goconst 块加一个 PluginName,值即 resultHandler.type
2

实现接口

新建 internal/plugin/event/<name>/,实现 types.WriterPluginExecute 里:taskCtx.ResultHandler 为 nil 或 config 解析失败时返回 *types.PluginError{Kind: types.PluginErrExecution, ...};用 types.VisitOutputRows(或 ConsumeOutputRows)遍历 outputs;写 IO 走传入的 io.Write;下游错误用 types.IsRetryableError(err) 决定 Retryable。尊重 taskCtx.DebugMode:只 PreviewOutputs,不发外部 IO。成功时返回要透传的数据(无则 nil)。
3

装配

internal/worker/app/app.gonewPluginRegistry 加 case;把名字加进 cmd/*/main.goPlugins 列表。若依赖新的 endpoint,在 validateProfileRuntimeConfig 里加启动期校验。
4

测试

包内单测(参考 blockdbwrite 的测试和 tableupserts/table_upserts_test.go:单行 / 多行 / 缺列 / 显式 nil / 写失败可重试性 / debug 模式);event/types/errors_test.goTestIsRetryableError 覆盖新错误分类;internal/worker/app/profile_test.goTestNewPluginRegistry_ProfileMembership;orchestrator 侧 TestOrchestrator_EventPhase_* 已覆盖通用路径,一般不用改。
Call.ArgsJSONNoResultCacheArgsPooled 都是进程内字段:Call 不能被 JSON 往返当作 args 的载体,ArgsPooled 只在恰好一个 call 引用该缓冲时才可置位。共享静态 args 的 builder 一律不设 NoResultCache

测试

  • internal/plugin/callbuilder/callbuilder_test.gointernal/plugin/event/event_test.go:注册表和”runner”式集成用例;internal/plugin/pipeline_integration_test.go:起 gRPC mock BlockDB 跑 DBScan → BlockDBWrite 全链路。
  • internal/plugin/callbuilder/types/operator_test.go:三值逻辑、大整数精确比较、legacy operator 转换、dedup key。
  • internal/plugin/callbuilder/bundlescan/stream_test.go:流式批边界、单行超限、PrepareStream 静态错误零 emit。
  • e2e:e2e/system/logical_type/dbscan_operator_test.go(真实 Worker 进程 + gRPC mock BlockDB),e2e/perf/s1_builder_dbscan_test.gos5_plugin_return_value_test.go。总览见 测试

相关文档

  • blockx 仓库 docs/specs/plugin-system.md:接口、数据结构、装配、DBScan / BlockDBWrite / BlockBundle / TableUpserts 详细设计(重点)。
  • blockx 仓库 docs/specs/architecture.md §4.2.3:Plugin 模块在整体架构中的位置。
  • blockx 仓库 docs/specs/substrait_ast.md:Operator(Plan)JSON 格式与 Substrait 差异。
  • blockx 仓库 docs/specs/2026-07-14-dbscan-in-memory-operator.md:DBScan 改用 GetBlockRows + 内存解释器的决策记录。
  • 站内:Worker(orchestrator 与阶段准入)、Call 执行子系统(Calls 阶段与流式派发)、IO 访问子系统TaskIOReader / TaskIO 的实现与限流)、Bundle 集群(bundlescan / bundlewrite 展开)、Task 生命周期