Skip to main content
Bundle 集群是 BlockX 面向**历史回填(bundle backfill)**的一套独立部署: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)、ReturnValue handler、devstub 等 block 能力——bundle profile 不注册它们,任务在准入时被拒;不改变 task 协议、slot 状态机、heartbeat 等公共语义(这些归 CoordinatorWorker);不编排上游路由(SDK 何时把流量切到 bundle endpoint 不在本组件内)。
  • 不变量 1:成员发现隔离。block worker 只注册到 workers 类 registry,bundle worker 只注册到 bundle-workers 类 registry;两个 Coordinator 各 watch 自己的精确前缀列表。Profile.Validate 在创建任何外部依赖前拒绝错配。
  • 不变量 2:同一套代码,不同 profileinternal/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 的 TuneDefaultsStreamBuild.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.Profileinternal/worker/app/app.go):Deployment / Service / WorkerRegistryPrefix / UsageService(chaintable-usage 服务名,由入口固定)/ Builders / Plugins / TuneDefaults。bundle 入口的实际值:
  • app.Profileinternal/coordinator/app/app.go):Deployment / Service / WorkerRegistryPrefixes []stringregisterCoordinatorServiceDeployment 决定注册哪个 gRPC service。
  • etcdapi.LegacyBundleWorkerRegistryPrefix = /blockx/bundle-workers/ProdDefaultBundleWorkerRegistryPrefix = /blockx/prod/lanes/default/bundle-workers/TestDefaultBundleWorkerRegistryPrefixapi/etcd/registry.go);ParseRegistryPrefix 返回 RegistryPrefix{Format, Environment, Lane, WorkerKind}
  • adapters.BundleCoordinatorServerinternal/coordinator/adapters/bundle_grpc_server.go):只做 message 转换,调 CoordinatorServer.ReserveWorkerSlot
  • bundlescan.BundleScanBuilderDecl / BundleSourceinternal/plugin/callbuilder/bundlescan/types.go):on-wire functionCallConfig.config,字段 blockBundle.numbertriggerSources[].tableId|bundleBucket|operator|function|code|param|forceResultCacheNoBundleFile = "NO_BUNDLE_FILE"
  • bundlescan.BundleScanBuilderbundlescan.go / stream.go):同时实现 types.CallBuildertypes.StreamingCallBuilderSetStreamLimits 在启动时注入批边界。
  • types.StreamingCallBuilder / CallStreamProducer / BatchEmit / ErrProductionStoppedinternal/plugin/callbuilder/types/types.go):流式两阶段契约。
  • bundlescan.DuckDBParquetReader / DuckDBParquetReaderConfigduckdb.go):StreamBundleWithPlan(ctx, objectURI, bundleNumber, plan, yield);filter 下推到 read_parquet(...) WHERE,dedup 在 Go 侧逐行做;S3 凭证懒初始化,ExpiredToken 且尚未 yield 时刷新 secret 重试一次。
  • bundlescan.BundleRowYieldparquet.go):借用行契约——行与其中的 string/[]byte 只在回调期间有效。
  • borrowed.Pool / borrowed.ScanResultinternal/duckdb/borrowed/):连接池与 chunk 迭代。
  • icebergsdk.ResolveDataFilesReq / IcebergIOAdapter / Planner.PlanDataFilesinternal/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 调 prepareBundleSourceresolveBundleSourceBucketbundleBucket 非空直接用;否则读 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_idvalidateDuckDBReadPlan + ToFilterSQLArgs 干跑 → 预建 ParquetJSONEncoder。最后 markProvablyUniqueSources 决定哪些 source 的 call 置 NoResultCache。任何一步失败即 BUILDER_FAILED,零 call 派发。
3

Calls 阶段:Run(scanSem 下,长)

Orchestrator 在 scanSem 下起 runCallProductionbundleScanProducer.Run 用 errgroup 限 3 并发扫 source。每行经 bundleRowCallBuilder.buildNext 编码成 Call(row-args 直接写进 argspool 缓冲区,ArgsPooled=true),bundleSourceBatchAppendermaxRowsmaxArgsBytes 先到者处 flushemitemit 阻塞在 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 前失败)→ InitWriteJobjob_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_idDebugMode 只预览输出,不建 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,保证 bundle deployment 只可能写 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.gopresigned_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.gointernal/worker/core/config.gointernal/worker/adapters/orchestrator.go;优先级:内置默认 → TuneDefaultsWORKER_CONFIG 文件 → 环境变量)。 bundle coordinator:COORDINATOR_REGISTRY_PREFIXES(逗号分隔精确列表)整体替换 profile 默认的 [/blockx/bundle-workers/, /blockx/prod/lanes/default/bundle-workers/];其余配置与 block coordinator 相同。

扩展点

  • 调 bundle profile 默认值:改 cmd/bundle_worker/main.goTuneDefaults;只许设默认,不许后处理最终配置(操作者显式值必须能覆盖)。对应更新 cmd/bundle_worker/profile_test.gointernal/worker/app/profile_test.go
  • 给 bundle worker 增/减 Builder 或 Plugin:改 profile 的 Builders / Plugins,并在 internal/worker/app/app.gonewBuilderRegistry / newPluginRegistry 里有对应 case;若新 plugin 也依赖 BatchWrite,把它加进 validateProfileRuntimeConfig
  • 改 bundlescan 的 source 语义(新字段、新的解析路径):bundlescan/types.go(decl)→ source_resolution.go(解析)→ bundlescan.goprepareBundleSource(所有静态校验必须在这里,两条路径共用)→ 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.goParseRegistryPrefix,两端 Validate 自动生效。
  • 新的 bundle coordinator RPC:改 api/grpc/bundlecoordinator/bundle_coordinator.protomake(见 Makefile proto 目标)→ internal/coordinator/adapters/bundle_grpc_server.go 只做转换,调度逻辑仍在共享 CoordinatorServer

测试

  • internal/plugin/callbuilder/bundlescan/*_test.gobundlescan_test.go(Build 路径、source 解析、archive 重映射)、stream_test.goTestBundleScanStream_*MaxRows / MaxArgsBytes 边界、单行超限、静态错误零 emit、sibling source 零 emit、emit 拒绝停扫、ctx 取消)、duckdb_test.go(reader 配置、句柄释放、连接池取消)、parquet_test.gojson_encoder_benchmark_test.go
  • internal/plugin/event/bundlewrite/bundle_write_test.gobatchwrite/*_test.go:Init/PUT/Commit/Cancel 状态机、错误分类、Content-Length 校验;tableupserts/table_upserts_test.go 覆盖两种模式。
  • cmd/bundle_worker/bundle_worker_process_test.goTestBundleWorkerProcess_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。
  • 测试组织总览见 测试
cmd/bundle-loadgen 用法见仓库内 cmd/bundle-loadgen/README.md--coordinator 指向 bundle coordinator,--input-uri 指向已验证唯一列的 parquet fixture,--workload noop|cpu|blockdb|function 选负载,--result-cache disabled|enabled--result-shape 做 CallCache / 输出编码的控制变量 A/B,--result-handler-type BlockBundleWriteResultHandler 可接真实 writer(必须用隔离测试表)。worker 侧必须 STREAM_BUILD_ENABLED=trueblockdb / function 场景依赖 worker 打开 STUB_BLOCKDB_DELAY_MS 等 stub。

相关文档

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 集群的采样方案。
站内:Plugin 系统(Builder / Writer 通用契约)、Worker(阶段准入、orchestrator)、Coordinator(调度、watcher)、Call 执行子系统IO 访问子系统(system cache、retry 分类)、后端适配器(BlockDB / Iceberg 客户端)、部署概览