深入解析 StarRocks BE 的 ConnectorBenchmark 模块:基于 Benchgen 的 Benchmark 连接器实现与模块边界约束
【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks
StarRocks BE 中的ConnectorBenchmark是连接器(Connector)体系内的一个特殊实现:它不读取任何真实的外部数据源,而是借助 Benchgen 数据生成引擎,在查询执行时按需合成 Schema 与数据,用于在不依赖存储、服务或完整 Exec 的情况下验证 Connector 契约与扫描链路。本文以 be/src/connector/benchmark/AGENTS.md 为骨架,结合其源码实现、构建配置与模块边界清单,讲解该模块的设计意图、类结构、数据流、可配置参数以及工程治理规则。
模块定位:为 Connector 契约提供"无外部依赖"的基准数据源
在 StarRocks 的 BE 端,Connector 抽象层承担了对外部数据源(Hive、Iceberg、JDBC、Elasticsearch、MySQL、文件系统、湖存储等)的统一读写接入。ConnectorBenchmark在其中扮演一个特殊角色:它是唯一一个不连接任何真实外部系统、而是由 Benchgen 按参数现场生成数据的连接器实现。其模块描述明确写道:
Benchgen-backed benchmark connector implementation above connector contracts without registry composition, storage, service, or full Exec coupling.
这句话概括了它的三重定位:
- Benchgen-backed:数据来源是
benchgen库,而非文件、网络或远端服务; - Above connector contracts:它只依赖
Connector、DataSourceProvider、DataSource这些连接器契约(位于 be/src/connector_primitive/connector.h); - 无 registry composition、无 storage、无 service、无完整 Exec 耦合:它刻意保持轻量,只做"扫数据、产出 Chunk"这一件事。
从代码结构看,该模块由四个文件组成,职责划分非常清晰:
| 文件 | 职责 |
|---|---|
| benchmark_connector.h / benchmark_connector.cpp | 实现BenchmarkConnector、BenchmarkDataSourceProvider、BenchmarkDataSource三层连接器主体 |
| benchmark_scanner.h / benchmark_scanner.cpp | 实现BenchmarkScanner,负责把 Benchgen 生成的 Arrow RecordBatch 转换为 StarRocks Chunk |
连接器三层结构:Connector → DataSourceProvider → DataSource
StarRocks 的 Connector 体系采用"连接器(Connector)→ 数据源提供者(DataSourceProvider)→ 数据源(DataSource)"的三级结构,ConnectorBenchmark完整实现了这一套契约:
BenchmarkConnector:继承starrocks::connector::Connector,connector_type()返回ConnectorType::BENCHMARK,通过create_data_source_provider()把扫描节点(ConnectorScanNode)和 Thrift 计划节点(TPlanNode)封装为BenchmarkDataSourceProvider。BenchmarkDataSourceProvider:继承DataSourceProvider,持有ConnectorScanNode*与TBenchmarkScanNode,负责按TScanRange创建BenchmarkDataSource;其中insert_local_exchange_operator()返回true(每个数据源独立算子、支持本地交换),accept_empty_scan_ranges()返回false(不接受空扫描范围)。BenchmarkDataSource:继承DataSource,是真正的数据生产单元,实现open()、get_next()、close()以及raw_rows_read()、num_rows_read()、num_bytes_read()、cpu_time_spent()等指标接口。
值得注意的是模块边界清单 be/module_boundary_manifest.json 中connectorbenchmark条目给出的约束:它只允许依赖ConnectorPrimitive、Expr、Runtime、ChunkCore、ColumnCore、Types、Common、Base、Gutil、StarRocksGen这些低层目标,而禁止引入connector/connector_registry.h与exec/exec_env.h。也就是说,这个模块只负责"实现",不负责"注册"——注册动作被明确要求放到ModuleBootstrap(be/src/module/connector_bootstrap.cpp 中实际 include 了connector/benchmark/benchmark_connector.h,与清单中modulebootstrap允许 includeconnector/benchmark/的规则完全一致)。
数据流剖析:从 TPlanNode 到 Chunk 的完整调用链
一个 benchmark 查询在 BE 端经历如下数据流:
- 参数解析(
BenchmarkDataSource::_init_params):从TBenchmarkScanNode读取db_name、table_name,从TBenchmarkScanRange读取start_row、row_count,组合成BenchmarkScannerParam(定义于 benchmark_scanner.h),并交给BenchmarkScanner。 - Benchgen 初始化(
BenchmarkScanner::open):把db_name通过benchgen::SuiteIdFromString()映射为SuiteId(若解析失败返回InvalidArgument("Unknown benchmark database: ...")),再调用benchgen::MakeRecordBatchIterator()创建RecordBatchIterator,拿到 ArrowSchema。 - 类型转换规划(
BenchmarkScanner::_init_converters):对每个 SlotDescriptor,用build_arrow_column_convert_plan()构建 Arrow→StarRocks 的转换函数树;若需要类型转换(need_cast为真),则通过VectorizedCastExprFactory::from_type()生成 Cast 表达式,否则直接用ColumnRef引用原始列。列缺失时返回NotFound("Benchmark column ... not found in schema")。 - 数据产出(
BenchmarkScanner::get_next):从迭代器拉取arrow::RecordBatch(_next_batch),按_max_chunk_size切片,用convert_arrow_array_to_column()逐列转换出 raw chunk,经raw_chunk->filter(_chunk_filter)过滤后,再逐列执行 Cast 表达式得到最终 Chunk。 - 外层循环(
BenchmarkDataSource::get_next):循环拉取直至 Chunk 非空或 EOF;EOF 时返回Status::EndOfFile,同时累加_rows_read与_bytes_read供指标统计。
这一链路与文件类连接器(如 be/src/connector/file 的 scanner)高度一致,唯一的差异在于"数据源"被替换成了 Benchgen 合成器,因此它非常适合用来在不搭建任何外部环境的前提下,验证 Connector 扫描算子、Arrow 类型转换、Cast 表达式与 Chunk 组装链路。
参数说明:db_name、table_name 与生成选项
BenchmarkScannerParam是连接器与 Benchgen 之间的唯一参数载体,其字段在 benchmark_scanner.h 中定义为:
struct BenchmarkScannerParam { std::string db_name; std::string table_name; benchgen::GeneratorOptions options; };各参数的来源与语义(结合_init_params与open的源码逻辑):
| 参数 | 来源 | 语义与默认值 |
|---|---|---|
db_name | TBenchmarkScanNode.db_name | 必须能通过SuiteIdFromString()映射为已知SuiteId,否则报InvalidArgument;等价于 Benchgen 中的 benchmark 套件(Suite)名 |
table_name | TBenchmarkScanNode.table_name | 套件内具体表名,用于在生成的 Arrow Schema 中定位字段;同时作为ArrowConvertContext.current_file参与转换上下文 |
options.scale_factor | TBenchmarkScanNode.scale_factor(可选) | 数据规模缩放因子,未设置时默认1.0 |
options.start_row | TBenchmarkScanRange.start_row | 从第几行开始生成,用于多扫描范围并行切分 |
options.row_count | TBenchmarkScanRange.row_count(可选);并受_read_limit约束 | 生成行数;扫描范围未携带时取-1(不限);若存在_read_limit,则取两者较小值 |
options.chunk_size | state->chunk_size() | 单次产出 Chunk 的行数上限;若非法则回退到state->chunk_size(),仍为 0 时兜底4096 |
这些字段来自 Thrift 定义(TBenchmarkScanNode、TBenchmarkScanRange,见 gensrc/thrift 目录下相关 thrift 文件),意味着扫描参数是由 FE 端在计划阶段下发的,BE 只负责消费。
构建与模块边界:如何编译、为什么这样隔离
ConnectorBenchmark的构建入口在 be/src/connector/CMakeLists.txt 第 309–332 行,由WITH_CONNECTOR_BENCHMARK开关控制:
if (WITH_CONNECTOR_BENCHMARK) ADD_BE_LIB(ConnectorBenchmark benchmark/benchmark_connector.cpp benchmark/benchmark_scanner.cpp ) target_link_libraries(ConnectorBenchmark PUBLIC ConnectorPrimitive Expr Runtime ChunkCore ColumnCore Types Common Base Gutil StarRocksGen) target_link_libraries(ConnectorBenchmark PRIVATE benchgen arrow) endif()从依赖拆分可以看出设计者的边界意图:
- PUBLIC 依赖全部是稳定的低层目标(连接器契约、表达式、运行时、列/块、类型、公共/基础工具、生成代码);
- PRIVATE 依赖只有
benchgen与arrow——这两个是 Benchmark 连接器独有的实现细节,被严格封闭在模块内部,不向外泄露。
同时,ConnectorBenchmark也被纳入 be/src/module/CMakeLists.txt 的ModuleBootstrap目标依赖中,与清单modulebootstrap条目允许的 include 前缀(含connector/benchmark/)吻合。整个 BE 模块体系(base、gutil、common、connectorprimitive、connectorbenchmark、modulebootstrap等)都在 be/module_boundary_manifest.json 中登记,每一条都声明了 owned roots、allowed include prefixes、allowed target deps 与 remediation 建议。
工程治理:AGENTS.md 的生成与机械校验
本文所依据的 be/src/connector/benchmark/AGENTS.md 本身是一份自动生成的模块边界文档(文件头明确标注BEGIN GENERATED: BE MODULE HARNESSES),它并非手写,而是由 be/module_boundary_manifest.json 渲染而来:
- 修改清单后执行
python3 build-support/render_be_agents.py --write可重新生成各模块的 AGENTS.md; - 执行
python3 build-support/check_be_module_boundaries.py --mode full可在 CI 中机械地验证同样的规则; - 对应的测试见 build-support/test_render_be_agents.py 与 build-support/test_check_be_module_boundaries.py。
对connectorbenchmark模块,清单给出的 Remediation 建议是:把该模块严格限定为 benchgen 连接器实现,注册动作放入 ModuleBootstrap,避免把 Connector 注册表、存储、服务或完整 Exec 代码拉入连接器库。这类"文档即规范、规范可校验"的做法,让模块边界不再停留在口头约定,而是变成可在构建期强制执行的工程约束。
小结
ConnectorBenchmark是理解 StarRocks BE Connector 架构的一个极佳切入点:它麻雀虽小,却完整覆盖了"连接器注册 → Provider 创建 → DataSource 扫描 → Benchgen 取数 → Arrow 转换 → Cast 表达式 → Chunk 产出"的整条链路,且不依赖任何外部系统,天然适合作为连接器契约的测试载体。通过阅读 benchmark_connector.cpp 与 benchmark_scanner.cpp 的源码,再对照 be/module_boundary_manifest.json 中的边界规则,读者既能掌握一个真实连接器的实现套路,也能理解 StarRocks BE 如何用机械校验来守住模块间的依赖纪律。
【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考