cmd/bundle_coordinator + cmd/bundle_worker。它复用与 block 集群完全相同的 Coordinator / Worker 装配代码,只在入口 profile 上收窄能力、切换 etcd registry 前缀并默认开启 stream build。系统全貌见 架构总览。
一个 block bundle 是按块高连续切分的一段数据(BundleRef.Number = N 覆盖高度 N*1000+1 … (N+1)*1000,bundle 0 只含高度 0),以 parquet 文件形式落在 S3 / Iceberg 表的 bundle=N 分区。bundle backfill 任务把这份 parquet 逐行喂给用户函数,再把全部输出用 BlockDB BatchWrite 一次性覆写目标表的同一个 bundle。
职责与边界
- 负责:接受
BlockBundleCallConfig/CallListCallConfig类型的 task;用 DuckDB 从 parquet 流式扫描行、按批产出 call;用BlockBundleWriteResultHandler/TableUpsertsResultHandler把输出经 BlockDB BatchWrite Job(Init → presigned PUT → Commit)写回。 - 不负责:单块实时 dbscan(
DBScanBuilder)、ReturnValuehandler、devstub 等 block 能力——bundle profile 不注册它们,任务在准入时被拒;不改变 task 协议、slot 状态机、heartbeat 等公共语义(这些归 Coordinator 与 Worker);不编排上游路由(SDK 何时把流量切到 bundle endpoint 不在本组件内)。 - 不变量 1:成员发现隔离。block worker 只注册到
workers类 registry,bundle worker 只注册到bundle-workers类 registry;两个 Coordinator 各 watch 自己的精确前缀列表。Profile.Validate在创建任何外部依赖前拒绝错配。 - 不变量 2:同一套代码,不同 profile。
internal/worker/app.Run/internal/coordinator/app.Run只有一份;cmd/*只声明Profile。bundle Coordinator 只暴露bundlecoordinator.v1.BundleCoordinatorService,block 只暴露coordinator.v1.CoordinatorService,proxy 按 gRPC method path 区分流量。 - 不变量 3:stream build 生产必开。bundle profile 的
TuneDefaults置StreamBuild.Enabled=true;操作者显式STREAM_BUILD_ENABLED=false仍可回滚到全量物化Build()路径。 - 不变量 4:静态错误零执行,动态错误部分执行。
PrepareStream在 Builder 阶段完成全部静态校验,失败则零 call 派发;Run中途失败时先前批次可能已执行。 - 不变量 5:BatchWrite 是唯一写路径,且启动期强校验。注册
BundleWrite(或 bundle profile 的TableUpserts)时blockdb.batchWriteAddr为空,进程在启动期失败,不等到 Writer 阶段。 - 不变量 6:Commit 结果不确定时不 Cancel、不建替代 Job,返回不可重调度失败并保留原
job_id。
代码位置
核心类型与接口
app.Profile(internal/worker/app/app.go):Deployment/Service/WorkerRegistryPrefix/UsageService(chaintable-usage 服务名,由入口固定)/Builders/Plugins/TuneDefaults。bundle 入口的实际值:
app.Profile(internal/coordinator/app/app.go):Deployment/Service/WorkerRegistryPrefixes []string;registerCoordinatorService按Deployment决定注册哪个 gRPC service。etcdapi.LegacyBundleWorkerRegistryPrefix=/blockx/bundle-workers/、ProdDefaultBundleWorkerRegistryPrefix=/blockx/prod/lanes/default/bundle-workers/、TestDefaultBundleWorkerRegistryPrefix(api/etcd/registry.go);ParseRegistryPrefix返回RegistryPrefix{Format, Environment, Lane, WorkerKind}。adapters.BundleCoordinatorServer(internal/coordinator/adapters/bundle_grpc_server.go):只做 message 转换,调CoordinatorServer.ReserveWorkerSlot。bundlescan.BundleScanBuilderDecl/BundleSource(internal/plugin/callbuilder/bundlescan/types.go):on-wirefunctionCallConfig.config,字段blockBundle.number、triggerSources[].tableId|bundleBucket|operator|function|code|param|forceResultCache;NoBundleFile = "NO_BUNDLE_FILE"。bundlescan.BundleScanBuilder(bundlescan.go/stream.go):同时实现types.CallBuilder与types.StreamingCallBuilder;SetStreamLimits在启动时注入批边界。types.StreamingCallBuilder/CallStreamProducer/BatchEmit/ErrProductionStopped(internal/plugin/callbuilder/types/types.go):流式两阶段契约。bundlescan.DuckDBParquetReader/DuckDBParquetReaderConfig(duckdb.go):StreamBundleWithPlan(ctx, objectURI, bundleNumber, plan, yield);filter 下推到read_parquet(...) WHERE,dedup 在 Go 侧逐行做;S3 凭证懒初始化,ExpiredToken且尚未 yield 时刷新 secret 重试一次。bundlescan.BundleRowYield(parquet.go):借用行契约——行与其中的 string/[]byte 只在回调期间有效。borrowed.Pool/borrowed.ScanResult(internal/duckdb/borrowed/):连接池与 chunk 迭代。icebergsdk.ResolveDataFilesReq/IcebergIOAdapter/Planner.PlanDataFiles(internal/sdk/iceberg/):一次解析整张表当前 snapshot 的全部 live 文件并编码为 bundle 索引;LookupBundleFileIndex(data, bundle)选文件。细节归 后端适配器。batchwrite.Writer.Submit(ctx, taskIO, initReq, outputs, columns) (Result, error)(internal/plugin/event/batchwrite/writer.go):计划 → Init → PUT → Commit;Result{JobID, RowCount, ContentBytes};commitStatusUnknownError永远Retryable() == false。bundlewrite.BundleWritePlugin/BundleWriteConfig{TargetTable, BlockBundle};tableupserts.NewPlugin(block)与tableupserts.NewBatchPlugin(bundle)。
数据流 / 执行流程
1
Reserve + Submit
调用方对 bundle coordinator 调
BundleCoordinatorService.ReserveWorkerSlot,拿到 worker_addr / slot_id 后向该 worker SubmitTask。调度、slot TTL、heartbeat 与 block 完全一致,见 Coordinator。2
Builder 阶段:PrepareStream(builderSem 下,短)
BundleScanBuilder.PrepareStream 反序列化 BundleScanBuilderDecl,对每个 source 调 prepareBundleSource:resolveBundleSourceBucket(bundleBucket 非空直接用;否则读 logical-types 判断 event/state,state 走 <tableId>._archive,再经 ResolveDataFilesReq 解析出唯一 parquet URI;NO_BUNDLE_FILE 或 snapshot 无该 bundle 则 skip、产 0 call)→ inline code 归一化 → archive state 表把 operator 里的 id 重映射为 original_id → validateDuckDBReadPlan + ToFilterSQLArgs 干跑 → 预建 ParquetJSONEncoder。最后 markProvablyUniqueSources 决定哪些 source 的 call 置 NoResultCache。任何一步失败即 BUILDER_FAILED,零 call 派发。3
Calls 阶段:Run(scanSem 下,长)
Orchestrator 在
scanSem 下起 runCallProduction,bundleScanProducer.Run 用 errgroup 限 3 并发扫 source。每行经 bundleRowCallBuilder.buildNext 编码成 Call(row-args 直接写进 argspool 缓冲区,ArgsPooled=true),bundleSourceBatchAppender 在 maxRows 或 maxArgsBytes 先到者处 flush → emit。emit 阻塞在 MaxOutstandingBatches 信用上,形成背压。单行 args 超过字节上限是 production error。emit 返回 ErrProductionStopped 表示 worker 受控停止(fast-fail),不算 builder 失败。4
Dispatch / Execute
通用路径,见 Call 执行子系统 与 Plugin 系统。
5
Writer 阶段:BatchWrite
Calls 收敛成功后,
BundleWritePlugin.Execute 校验 targetTable / blockBundle,读列定义,调 batchwrite.Writer.Submit:先无保留遍历 outputs 算精确行数与字节数(planDelimitedWriteRows;超过 10,000,000 行或 5,000,000,000 bytes 在 Init 前失败)→ InitWriteJob 拿 job_id + data_url → 边编码边通过 io.Pipe 单次 HTTP PUT(准确 Content-Length,不跟随重定向)→ CommitWriteJob(row_count)。PUT 失败或 Commit 前 ctx 取消则 best-effort CancelWriteJob(5s 独立超时,context.WithoutCancel)。成功返回 written_rows / finish_time / job_id。DebugMode 只预览输出,不建 Job。TableUpsertsResultHandler 在 bundle worker 上走同一个 batchwrite.Writer,只是策略换成 NormalTableUpsertStrategy{Condition, UpdateColumns};planNormalUpsertDelimitedWriteRows 在 Init 前要求每行有效字段集合一致、id 非空、至少一个非 ID 列,并把首行推导出的 update_columns 显式写入请求。block worker 上的同名 handler 仍调 TableWriteService.UpsertRows(sync=false)。写路径由 profile 在装配期固定,task 配置里没有 writeApi。状态与生命周期
- registry 归属:worker 启动时
Profile.Validate+ 配置加载后再次validateWorkerRegistryPrefix,保证bundledeployment 只可能写bundle-workers类前缀(legacy 或 v2)。coordinator 端validateRegistryPrefixes同理,并拒绝重复前缀。 - stream build 两阶段:Builder 阶段只持有 builderSem 做
PrepareStream;Calls 阶段持有 scanSem + executorSem 跑 production;production 结束或失败后释放 scanSem。一个 scanSem 是一个 task 的生产,内部最多 3 个 source 并发,故 worker 级 DuckDB 扫描并发上限约为ScanSlots × 3。 - BatchWrite 失败分类(
batchwrite/classified_error.go、presigned_upload.go):配置/校验/编码/非法响应/重定向/其他 4xx → 不可重调度;Meta/BlockDB 瞬时 IO、HTTP 传输错误、PUT 408/425/429/500/502/503/504 → 可重调度;Commit 请求或响应不确定 →commitStatusUnknownError,不可重调度且保留job_id。插件把分类透传为PluginError.Retryable。 - Iceberg 解析缓存:
ResolveDataFilesReq.CacheKey()只用逻辑表 ID,走 Worker system cache(默认 TTL 20 分钟,SystemIOCacheTTLms);AcceptCachedRead在缓存索引的高水位覆盖RequiredBundle时复用,否则整表刷新一次。
配置
以下是 bundle worker 最相关的配置(结构体字段来自internal/worker/app/config.go、internal/worker/core/config.go、internal/worker/adapters/orchestrator.go;优先级:内置默认 → TuneDefaults → WORKER_CONFIG 文件 → 环境变量)。
bundle coordinator:
COORDINATOR_REGISTRY_PREFIXES(逗号分隔精确列表)整体替换 profile 默认的 [/blockx/bundle-workers/, /blockx/prod/lanes/default/bundle-workers/];其余配置与 block coordinator 相同。
扩展点
- 调 bundle profile 默认值:改
cmd/bundle_worker/main.go的TuneDefaults;只许设默认,不许后处理最终配置(操作者显式值必须能覆盖)。对应更新cmd/bundle_worker/profile_test.go、internal/worker/app/profile_test.go。 - 给 bundle worker 增/减 Builder 或 Plugin:改 profile 的
Builders/Plugins,并在internal/worker/app/app.go的newBuilderRegistry/newPluginRegistry里有对应 case;若新 plugin 也依赖 BatchWrite,把它加进validateProfileRuntimeConfig。 - 改 bundlescan 的 source 语义(新字段、新的解析路径):
bundlescan/types.go(decl)→source_resolution.go(解析)→bundlescan.go的prepareBundleSource(所有静态校验必须在这里,两条路径共用)→source_calls.go(call 拼装)。不要把校验放进Run,否则破坏静态失败零执行。 - 改 DuckDB 读取:SQL 拼装、S3 secret、filter 下推、dedup 在
bundlescan/duckdb.go;连接池 / 行解码 / 类型支持在internal/duckdb/borrowed/(不放业务策略)。 - 改 BatchWrite 生命周期(重试、上传方式、容量上限):
internal/plugin/event/batchwrite/。BlockDB SDK 只保留InitWriteJobReq/CommitWriteJobReq/CancelWriteJobReq三个原子请求,不要把编排下沉到 SDK。 - 新增一种 registry 前缀格式:
api/etcd/registry.go的ParseRegistryPrefix,两端Validate自动生效。 - 新的 bundle coordinator RPC:改
api/grpc/bundlecoordinator/bundle_coordinator.proto→make(见Makefileproto 目标)→internal/coordinator/adapters/bundle_grpc_server.go只做转换,调度逻辑仍在共享CoordinatorServer。
测试
internal/plugin/callbuilder/bundlescan/*_test.go:bundlescan_test.go(Build 路径、source 解析、archive 重映射)、stream_test.go(TestBundleScanStream_*:MaxRows/MaxArgsBytes边界、单行超限、静态错误零 emit、sibling source 零 emit、emit 拒绝停扫、ctx 取消)、duckdb_test.go(reader 配置、句柄释放、连接池取消)、parquet_test.go、json_encoder_benchmark_test.go。internal/plugin/event/bundlewrite/bundle_write_test.go与batchwrite/*_test.go:Init/PUT/Commit/Cancel 状态机、错误分类、Content-Length 校验;tableupserts/table_upserts_test.go覆盖两种模式。cmd/bundle_worker/bundle_worker_process_test.go:TestBundleWorkerProcess_RejectsMissingBatchWriteAddr(缺BLOCKDB_BATCH_WRITE_ADDR必须启动失败)、TestBundleWorkerProcess_BootsWithBundleProfile(启动、就绪、只注册 bundle profile 能力);profile_test.go钉住 deployment/service/registry 值。cmd/bundle_coordinator/*_test.go同理。internal/coordinator/adapters/bundle_grpc_server_test.go:bundle adapter 复用共享调度并透传 gRPC status。- 测试组织总览见 测试。
相关文档
blockx 仓库docs/specs/:
docs/specs/2026-07-21-block-bundle-clusters.md— block / bundle 部署拆分:双入口、双 registry、系统不变量(重点)。docs/specs/2026-07-28-registry-prefix-migration.md— legacy → v2 registry 前缀格式、WORKER_REGISTRY_PREFIX/COORDINATOR_REGISTRY_PREFIXES覆盖规则。docs/specs/2026-07-23-bundle-coordinator-api.md— 独立BundleCoordinatorService的决策与上线顺序。docs/specs/2026-07-30-bundle-write-batch-api.md— BatchWrite 写路径的执行流程、契约边界、失败状态机(规范性)。docs/specs/plugin-system.md§6.1(StreamingCallBuilder契约、NoResultCache四条件)、§10(BlockBundleCallConfig/BlockBundleWriteResultHandler/TableUpsertsResultHandler契约)。docs/specs/2026-07-30-ec2-worker-profiling.md— 面向 block / bundle 两套 EC2 Worker 集群的采样方案。