Skip to main content
后端适配器是 BlockX 访问外部服务的客户端层,代码全部在 internal/sdk/(目录名沿用历史,里面是各服务的 client,不是对外发布的 SDK)。每个子包把一个远端服务(BlockDB gRPC、链节点 JSON-RPC、Meta gRPC、Iceberg/Glue)包装成 IO 子系统认识的 iocore.BackendAdapter,再由 Worker 或 sync-invoker 在启动时装配进 WorkerIOScope。Builder、Writer Plugin 和 Function Code View 通过 TaskIO.Read/Write 或直接持有 adapter 使用它们。整体位置见 架构总览,调用链的上半段见 IO 访问子系统

职责与边界

  • 每个 SDK 包只负责一件事:把 iocore.ReadReq / iocore.WriteReq 翻译成一次远端调用,并把响应放进 ReadResult.Data / WriteResult.Data。请求类型由 SDK 包定义,不由 internal/io/core 定义。
  • SDK 不做缓存、不做 singleflight、不做 task 窗口、不做准入控制。这些都在 IO 子系统;SDK 只提供 Registration(...)assembly.Module 帮助调用方把 adapter 包进 adaptive.WrapBackend
  • SDK 负责错误分类:实现 iocore.ErrorClassifier.ClassifyError(决定 retryable / non-retryable / param)和 adaptive.OutcomeClassifier.ClassifyAdaptiveOutcome(决定 AIMD 是否收缩)。本地 context 取消/超时一律不算下游过载。
  • BlockDB SDK 只暴露原子 RPC。BatchWrite 的 Init → PUT → Commit 编排在 internal/plugin/event/batchwrite,不在 SDK(docs/specs/2026-07-30-bundle-write-batch-api.md §2)。
  • 面向 Python executor 的 bridge 路径(BlockDB bridge、NodeRPC bridge、capability backend)只透传 protobuf/JSON-RPC 字节,不重建请求体。
  • 连接策略属于 adapter 实例:指向 NLB 的 gRPC client 用 grpcclient.NewMultiConn,每个 adapter 内相同地址复用一个 ClientConn,没有进程级连接注册表。
  • internal/duckdb/borrowed 只提供 DuckDB 原生绑定之上的连接池与结果解码;SQL 拼装、S3 凭证、过滤下推留在 internal/plugin/callbuilder/bundlescan

代码位置

核心类型与接口

所有 adapter 都实现同一个接口(internal/io/core/model/adapter.go):

blockdb

  • AdapterConfigadapter.go):每个 gRPC service 一个地址字段(TableReadAddrTableWriteAddrTableScanAddrBlockReadAddrBlockWriteAddrBatchWriteAddrBlockSubscribeAddrL2BlockAddrTimeReadAddrTimeWriteAddr)加 MaxRecvMsgSize(0 表示 32 MiB)。相同地址复用一个 ClientConn;未配置的 service 在被调用时返回 blockdb:service_not_configured(param 类错误)。TableSubscribeAddr 字段存在并从 BLOCKDB_TABLE_SUBSCRIBE_ADDR 读入,但 NewAdapter 不会拨号它;它只在 internal/worker/app/function_code.go 里作为启用 BlockDB Function Code View 的条件之一。
  • Adapteradapter.go):typed 后端,BackendBlockDB = "blockdb"Read 用类型开关分派 ReadOp[*blockdbpb.GetRowRequest]BatchGetRowsRequestFilterRowsRequestGetBlockRowsRequestGetEventRequestGetStateRequestGetValueRequest,把 ResultSet 转成 JSON 行数组;JSON/LIST/DICT 类逻辑列按 ColumnMeta.logical_type 还原为对象。Write 分派 WriteOp[*UpsertRowsRequest]WriteOp[*DeleteRowsRequest]WriteOp[*UpsertTimeRowsRequest]BlockWriteReqInitWriteJobReqCommitWriteJobReqCancelWriteJobReq
  • ReadOp[T] / WriteOp[T]types.go):泛型包装,把任意 proto 请求变成 iocore.ReadReq/WriteReq,附带 Key(cache key)、Op(统计名)、RTO/WTO(超时)。FilterEqFilterAndFilterFilterJSON 生成参数化 WHERE;Value(any)*blockdbpb.Value
  • 三种表语义(proto/table.protoblock.prototime.proto):Table(L1 行表:GetRow/BatchGetRows/FilterRows/UpsertRows/DeleteRows/Scan)、Block(L2 区块作用域:GetEvent/GetStatecurrent_blockUpsertBlockBlock{block_id,height,timestamp}L2BlockReadService.GetBlockRows)、Time(GetValue(row_id,time_at)UpsertTimeRows(time_at))。BlockWriteReq 是多表 L2 写:Data map[table][]*WriteRow,adapter 对每张表发一次 UpsertBlock
  • BatchWrite(batch_write.goproto/batch.proto):InitWriteJobReq/CommitWriteJobReq/CancelWriteJobReq 各包一个 BatchWriteService 请求,响应以 proto 字节放进 WriteResult.Data。上传 presigned URL、校验 job_id 一致性等编排在 internal/plugin/event/batchwrite/writer.go
  • BridgeAdapterbridge_adapter.go):BackendBlockDBBridge = "blockdb_executor",接收 Python blockdb_bridge.BridgeChannel 经 UDS 二进制 sidecar 发来的原始请求 protobuf,用 capabilitywire.DecodeBinaryRequest 解出 gRPC method path 后按 bridgeReadMethods / bridgeWriteMethods 白名单转发;bridgeBlockedWriteMethods 拒绝函数代码发起 BatchWrite 生命周期;Scan/Subscribe 流式方法不支持。响应字节原样放进 ReadResult.Data
  • ScanClient / Adapter.ScanAllonline.go):TableScanService.Scan 的流式读取;Function Code View 复用 typed Adapter.ScanAllSubscribeClient 目前只有构造与连接持有,没有 Subscribe 方法。
  • 错误分类(adapter.goadaptive.go):classifyBlockDBError 解析 gRPC status 里的 Error NNNN (SQLSTATE) 得到 mysql:NNNN / proxysql:NNNN code,再按 gRPC code 映射;classifyAdaptiveOutcomeResourceExhausted/DeadlineExceeded/Unavailable、传输错误和一组 MySQL/TiDB/ProxySQL 超时短语判为过载。
  • Registration(cfg, adaptive.Config, opts...) / BridgeRegistration(...):返回 iocore.BackendRegistrationWithMultiConnConfig 注入 grpcclient.MultiConnConfig
  • 测试替身:DelayedGRPCStubgrpc_stub.go,供 cmd/blockdb_grpc_stub 挂到真实 gRPC server)、DelayedBridgeStubAdapterbridge_stub.go,进程内 bridge stub,STUB_BLOCKDB_DELAY_MS 触发)。

noderpc

  • NodeRPCReqtypes.go):Go 侧直接调用的 JSON-RPC 2.0 请求,BackendNodeRPC = "rpc"Operation() 返回 method,Timeout() 固定 10s。
  • Adapteradapter.go):Read 把请求体 POST 到 endpoint;bridge 请求(Python executor 发来的 rawBridgeRequester)从 JSON 元数据取 params.chain,POST 到 endpoint/<chain>,body 直接使用 sidecar 字节。Write 不支持。每个请求带 x-load-deadline/x-load-priority/x-load-retries 头,HTTP 2xx 中的 {"error":...} 解析为 jsonRPCError,非 2xx 变成 HTTPStatusError
  • HTTPClientConfig{TransportCount, RefreshIntervalMs} + rotatingHTTPTransporttransport.go):8 个独立 http.Transport shard 轮询,每 30s 平滑替换一个;HTTPPoolSizesForAdmission 从 admission 初始/上限并发推导 idle 连接预算。
  • Registration(endpoint, adaptive.Config, opts...)
  • 语义批处理(batching_backend.gobatch_core.gobatch_protocol.gobatch_config.gobatch_window*.go):BatchingBackend 实现 BackendAdapter 并包在 admission 之外(assembly.Module.OuterWrap)。BatchConfig{Mode, MaxWait, MaxItems, FallbackCooldown}BatchMode = off|shadow|onBatchConfigFromEnv()NewBatchingBackend(inner, cfg, BatchObserver)。可批方法只有 getAddressCodegetStorageAtgetAddressBalancegetAddressNonce,且 block context 必须是固定 hash/height;物理请求是内部方法 blockx_stateReadBatch,载荷是 BSRB/1 二进制。batchobserved.Observer() 提供 obs 指标钩子。
  • 二进制 sidecar:bridge 请求的直通路径在 Adapter.readFromBridge——缺少 RawBinaryRequest() 直接返回 noderpc:invalid_bridge_request,没有 JSON 重建回退。Python 侧在 python/blockx_executor/noderpc_bridge.py

meta 与 logicaltypes

  • internal/sdk/meta/proto/*.protogo_package 指向 blockdb 仓库;make proto-meta--go_opt=M... 把三个文件重映射到 internal/sdk/meta/gen(包名 gen,调用方通常别名 metapb)。
  • logicaltypes.Adapterlogicaltypes/adapter.go):BackendLogicalTypes = "logicalTypes"ReadReq{TableID, RTimeout}metapb.TableMetaServiceClient.GetTable → JSON 编码的 TableSchema{Type string, Columns []logical_types.Column}ClassifyError 把 schema/配置类 gRPC code 判为不可重试。LocalStubAdapter 按 table ID("*" 兜底)返回本地 schema。Registration(host, adaptive.Config, opts...)
  • meta.Adaptermeta/meta.go):HTTP GET /api/v1/meta/get_table_metadata?id= 的旧客户端,BackendMeta = "meta"MetaReadReq。仓库内没有生产代码 import 它,只有包内单测和 //go:build livemeta_live_test.go。Worker 的 META_ADDR 走的是 logicaltypes

iceberg

  • Config{Namespace, Region}NewPlanner(ctx, cfg)planner.go):用 AWS 默认凭证链(EC2 上即 instance role)构造 Glue catalog,构造时不访问 Glue/S3。DataFilePlanner 接口只有 PlanDataFiles(ctx, tableName) ([]string, error)Planner.PlanDataFiles 加载表、规划当前快照全部 data file,遇到 delete file 直接报错。
  • ResolveDataFilesReq{TableID, PhysicalTableID, RequiredBundle, RTimeout}io_adapter.go):BackendIcebergResolve = "iceberg"CacheKey() 只用逻辑 TableID;实现 CachedReadValidator.AcceptCachedRead——缓存索引的 high water ≥ RequiredBundle 才复用,否则触发一次全表刷新(对应提交 “iceberg: refresh indexes above bundle high water”)。
  • IcebergIOAdapter:只读,把路径列表编码成 BFI1 索引(bundle_file_index.go:公共前缀 + 每 bundle 一条记录,同一 bundle 出现两个 live 文件即报错)。LookupBundleFileIndex(encoded, bundle)BundleFileIndexHighWater(encoded)bundlescan 解码。
  • Worker 把它注册为 SystemCache: trueReadOnly: true 的模块,准入是固定并发 IO.SystemIOCacheRefreshConcurrency(默认 2)。所需 IAM 见 blockx 仓库 docs/deploy.md §3.3(glue:GetTables3:GetObject,SSE-KMS 时加 kms:Decrypt)。

localtestservice

  • Serviceservice.go)直接实现生成的 LocalTestServiceServerHelloWorld 返回 "hello, <name>"BlockDBGetRow 通过注入的 BlockDBReader(宿主的 raw blockdb adapter 或 stub)发一次 ReadOp[*blockdbpb.GetRowRequest]
  • Adapteradapter.go):BackendLocalTestService = "localtestservice",拆 capabilitywire 二进制信封、按 method path proto.Unmarshal、进程内调 ServiceModule(blockDB BlockDBReader) assembly.Module:只读、固定准入 4。
  • examples/local_test_service/ 是 Python 侧端到端示例(blockx-py 的 LocalTestService 类 → BridgeChannel → worker);新增 capability backend 的逐文件模板见 blockx 仓库 docs/capability-backend-guide.md

router

  • Adapterrouter/adapter.go):BackendRouter = "router",只支持 routerpkg.RouterFindFunctionsFullMethodName,业务在依赖 github.com/chaintable/router/goService.FindFunctionsModule() 只读、固定准入 100。

grpcclient

  • New(target, opts...):DNS round_robin 的 insecure ClientConn。
  • NewMultiConn(target, MultiConnConfig, opts...):注册名为 blockx_multi_conn 的自定义 balancer,在一个 ClientConn 内维持 Connections 条 SubConn,每 RefreshIntervalMs 换一条;DefaultMultiConnConfig() = 8 条 / 10000 ms。grpc.WithDisableServiceConfig() 防止 resolver 下发的 service config 覆盖策略。
  • DelimitedWriteRowEncoderdelimited_write_row.go):按 logical_types.Column 决定 Value oneof kind,直接写 length-delimited blockdb.v1.WriteRow 字节,不构造生成的 proto message;BatchWrite 上传文件用它。

duckdb/borrowed

  • Poolpool.go):NewPool(db, maxOpen, maxIdle)Acquire(ctx)Release(conn)Close()WatchContext(ctx, conn)context.go)用 duckdb_interrupt 把 context 取消传给 DuckDB;ScanResult(ctx, res, yield)rows.go)按 chunk 迭代并把 VARCHAR/BLOB 以零拷贝视图交给回调——回调返回后值失效,需要保留必须复制。
  • 转换逻辑改编自 duckdb-go,直接调用 duckdb-go-bindingsLICENSE.duckdb-goLICENSE.duckdb-go-bindingsREADME.md 记录来源与所有权。

cmd/blockdb_grpc_stub

  • main.goblockdb.RegisterDelayedGRPCStub(registrar, readDelay) 起一个真实 gRPC server:读 RPC 等待 -read-delay 后返回固定 token 形状(GetRowGetState),写 RPC 立即返回;-probe-address 模式只探测已运行的 stub 并校验其 delay。用途和 worker 侧环境要求见 cmd/blockdb_grpc_stub/README.md

数据流 / 执行流程

1

装配

internal/worker/app/app.go(sync-invoker 在 cmd/syncinvoker/io.go)把每个后端声明为一个 assembly.Module{Kind, Enabled, Build, Stub, Admission, OuterWrap}Enabled 为假时用 Stub(如 NODE_RPC_ENDPOINT 为空用 devstub,META_ADDR 为空用 logicaltypes.LocalStubAdapter)。assembly.Build 统一经 adaptiveobserved.WrapBackendRegistration 套上带指标的准入包装(internal/io/assembly/assembly.go)。
2

请求进入

Go 调用方构造 SDK 请求类型(如 blockdb.ReadOp[*blockdbpb.GetBlockRowsRequest])交给 TaskIO.Read;Python 调用方的 gRPC/JSON-RPC 字节经 UDS sidecar 到 Worker,被包成带 RawBinaryRequest() 的 IORequest,Backend() 由 wire 上的 backend 名决定。
3

SDK 执行

adapter 只看自己认识的请求类型;不认识的返回错误。gRPC 调用前用 clientid.IntoOutgoingGRPC(ctx) 传播 client id,HTTP 用 clientid.IntoHTTP
4

错误回流

adapter 返回的 error 先经 ClassifyAdaptiveOutcome 喂给 AIMD,再经 ClassifyError 变成 IOError 决定 Call/IO 层重试。两处都把本地 context 结束判为 Ignore / timeout,而不是下游过载。

状态与生命周期

  • gRPC 连接:NewMultiConn 的 balancer 为每个 slot 先建 candidate SubConn,READY 后才提升并排空旧连接;首次刷新在 (0, refreshInterval] 内随机延迟。Adapter.Close() 关闭全部 ClientConn
  • NodeRPC transport:rotatingHTTPTransport 每个 refresh 间隔退役一个 shard,退役 shard 等在飞请求结束再关闭;Close() 拒绝新请求并等待全部排空。
  • 批处理 group:groupCollecting → groupFlushing → groupDone;同 chain、同固定 block、同 timeout、同 client id 的请求进同一 group,MaxWait(默认 500µs,Linux 上用 timerfd 计时)或 MaxItems(默认 16)触发 flush。收到 -32601 进入 Worker 级 cooldown(默认 60s),到期后只放一个 half-open 探针 batch。Close() 让收集中的 group 失败并取消在飞物理请求。
  • 幂等:BlockDB 各 RPC 幂等,所以 IO 层可以对 Result 阶段的 Init/Commit/Cancel 做有限重试;NodeRPC 只有 Read 路径。

配置

Worker 的字段在 internal/worker/app/config.goWorkerFullConfig),环境变量覆盖在 internal/worker/app/app.go;默认值从 DefaultWorkerFullConfig() 读。

扩展点

  • 新增一个 BlockDB RPC:改 internal/sdk/blockdb/proto/*.protomake proto-blockdb → 在 Adapter.dispatchRead/dispatchWriteReadOp[*NewRequest] 分支 → 如需对 Python 开放,把 FullMethodName 加进 bridge_adapter.gobridgeReadMethods / bridgeWriteMethods → 在 blockdb_test.go(bufconn + gomock server)加用例。
  • 新增一个 backend:实现 BackendAdapter(可选 ErrorClassifierOutcomeClassifier),提供 RegistrationModule,然后在 internal/worker/app/app.gobackendModules 加一条声明。gRPC 风格的 capability backend 直接照 internal/sdk/localtestservice 复制,步骤在 docs/capability-backend-guide.md
  • 新增可批处理的 NodeRPC 方法:改 batch_core.gobatchMethodsbatch_protocol.go 的 BSRB kind,同时改 Leafage 侧 handler,并遵循 docs/specs/noderpc-semantic-batching.md §7.3 的版本演进规则。
  • 调整连接策略:只改对应 adapter 的 grpcclient.MultiConnConfig,不要新增全局连接池。
  • 重新生成 proto:需要 protocprotoc-gen-goprotoc-gen-go-grpc$(go env GOPATH)/binMakefile 顶部有安装命令)。make proto 会先跑 proto-blockdbproto-metaproto-localtestservice,再生成 api/grpc/*。三个 SDK 目标都带 require_unimplemented_servers=false

测试

  • blockdbblockdb_test.gobufconn 起 gRPC server,mock_server_test.go 是 mockgen 生成的 service mock;partial_config_test.go 覆盖未配置 service 的干净失败;batch_write_test.go 校验 proto 与 BlockDB 仓库契约一致;bridge_*_test.go 覆盖 bridge 白名单与原始字节路径。
  • noderpcnoderpc_test.gohttptest.Serverbatching_backend_test.go 覆盖跨 task 合批、错误隔离、-32601 回退与 cooldown 探针;batch_window_linux_test.go 在 timerfd 不可用时 t.Skip
  • 需要外部服务的测试都会跳过:TestAdapter_ReadReal 需要 NODERPC_REAL_URL(且非 -short);meta_live_test.go 需要 -tags liveMETA_API_URL。CI 的 make test-short,因此不会碰真实服务。
  • icebergplanner_test.go 只测 delete-file 拒绝和索引编码,不访问 Glue。
  • duckdb/borrowedpool_test.gocells_test.gorows_test.go 需要 cgo 与 DuckDB 绑定,随 go test 编译。
  • E2E:e2e/system/logical_type/metapb 起本地 Meta gRPC 服务测 logical types 与 dbscan;进程级测试各自起 mock gRPC server(cmd/worker/worker_process_blockdb_bridge_test.go 实现 TableReadServiceServercmd/worker/worker_process_function_code_test.gointernal/testutil.StartMockFunctionBlockDB);性能基准用 cmd/blockdb_grpc_stub

相关文档

  • blockx 仓库 docs/specs/2026-07-30-bundle-write-batch-api.md:BatchWrite Job 生命周期,SDK 只提供三个原子 RPC。
  • blockx 仓库 docs/specs/noderpc-semantic-batching.md:批处理 facade 的准入规则、group 状态机、BSRB/1 协议、配置与回退。
  • blockx 仓库 docs/specs/2026-08-06-noderpc-binary-sidecar.md:NodeRPC bridge 二进制直通。
  • blockx 仓库 docs/specs/2026-07-23-grpc-client-balancing.mdNewMultiConn 的设计与各组件配置。
  • blockx 仓库 docs/specs/io-subsystem.mddocs/io-backend-module-design.mddocs/capability-backend-guide.md:backend adapter 契约与新增后端指南。
  • blockx 仓库 docs/deploy.md §3.3(Iceberg IAM)、§9.2(Worker 环境变量)。
  • 站内:IO 访问子系统Plugin 系统Function Code ViewPython ExecutorBundle 集群协议与接口测试组织与命令