尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

StarRocks BE ConnectorLake 模块:共享数据湖 Lake 连接器与延迟物化实现解析

发布时间:2026/9/15 21:42:29

资讯中心
01
ARTICLE

StarRocks BE ConnectorLake 模块:共享数据湖 Lake 连接器与延迟物化实现解析

StarRocks BE ConnectorLake 模块:共享数据湖 Lake 连接器与延迟物化实现解析
StarRocks BE ConnectorLake 模块共享数据湖 Lake 连接器与延迟物化实现解析【免费下载链接】starrocksThe worlds 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中be/src/connector/lake模块展开聚焦共享数据shared-data湖表扫描连接器LakeConnector、LakeDataSource以及湖表全局延迟物化上下文LakeScanLazyMaterializationContext的实现与模块边界。读完本文你将掌握 Lake 扫描链路的类职责、关键参数与谓词下推路径理解该模块如何在模块边界清单约束下保持与 Storage 及底层计算/运行时常量的解耦并能基于源码定位后续深入阅读的入口。模块定位Concrete Shared-Data Lake Connector在 StarRocks BE 的模块化改造中ConnectorLake内部标识connectorlake被明确定义为一个位于 Storage 之上的具体共享数据湖连接器模块。其核心职责记录在 be/src/connector/lake/AGENTS.md 中Concrete shared-data lake connector and lake late-materialization context above Storage, without registry composition, service, or full Exec coupling.翻译过来即提供共享数据湖表的扫描连接器以及在 Storage 之上维护湖表延迟物化状态同时不依赖注册表组合registry composition、服务层service或完整的 Exec 耦合。该模块拥有两个顶层构建目标ConnectorLake和测试目标connector_lake_test其职责被进一步细化为只负责 shared-data 湖表扫描与湖表延迟物化状态管理将注册动作放在 ModuleBootstrap即 be/src/module/connector_bootstrap.cpp中完成避免与注册表connector registry、服务层和完整 Exec 层耦合。从代码结构看be/src/connector/lake目录下仅有 4 个源文件正好覆盖上述两项核心职责文件职责lake_connector.hLakeConnector与LakeDataSource、LakeDataSourceProvider的声明lake_connector.cppLake 扫描连接器的完整实现约 2000 行lake_global_late_materialization_context.h湖表延迟物化上下文LakeScanLazyMaterializationContext的声明lake_global_late_materialization_context.cpp延迟物化上下文的实现模块边界契约允许的依赖与禁止的耦合ConnectorLake 的边界不是口头约定而是由 be/module_boundary_manifest.json 中的结构化条目id为connectorlake见第 1116–1150 行机械地约束并可通过构建脚本校验。允许的 include 前缀connector/lake/、connector_primitive/、storage/、storage_primitive/、compute_env/、exec_primitive/、exprs/、runtime/、platform/、fs/、io/、column/、types/、common/、base/、gutil/、gen_cpp/。允许的目标依赖allowed target depsConnectorPrimitive、Storage、ComputeEnv、StoragePrimitive、ExecPrimitive、Expr、Runtime、Platform、FileSystem、IO、ChunkCore、ColumnCore、Types、Common、Base、Gutil、StarRocksGen。明确禁止的耦合禁止的 include 前缀connector/模块内connector/lake/除外、exec/、service/、http/、agent/、script/、tools/、formats/、cache/、util/禁止包含的头文件connector/connector_registry.h与exec/exec_env.h。这意味着 ConnectorLake 可以自由地站在 Storage 之上读取湖表数据、复用底层谓词树storage_primitive、全局字典compute_env与运行时常量runtime但不能向上触碰连接器注册表、ExecEnv 单例、HTTP 管理与 agent 等高层设施。这种边界设计的价值在于当开发者想为 Lake 扫描引入一个新功能时清单会直接给出往哪里放、往哪里不许放的明确指引。该模块边界同时由两个 Python 脚本保障修改清单后运行python3 build-support/render_be_agents.py --write重新生成各模块的 AGENTS.mdbe/src/connector/lake/AGENTS.md顶部的BEGIN GENERATED标记即由此产生运行python3 build-support/check_be_module_boundaries.py --mode full可机械化校验同样的规则。LakeConnector 与 LakeDataSource湖表扫描的执行骨架Connector 接口与注册LakeConnector是Connector抽象类的具体实现connector_type()返回ConnectorType::LAKE。它与 Hive、File、Iceberg 等连接器一同通过 be/src/module/connector_bootstrap.cpp 中的bootstrap_builtin_connectors()注册install_if_absentHiveConnector(registry, Connector::HIVE); install_if_absentFileConnector(registry, Connector::FILE); install_if_absentLakeConnector(registry, Connector::LAKE); install_if_absentCacheStatsConnector(registry, Connector::CACHE_STATS);install_if_absent采用若注册表中不存在同名连接器才安装的幂等策略避免重复注册。LakeConnector唯一覆写的方法create_data_source_provider()会基于TPlanNode构造一个LakeDataSourceProviderDataSourceProviderPtr LakeConnector::create_data_source_provider(ConnectorScanNode* /*scan_node*/, const TPlanNode plan_node) const { return std::make_uniqueLakeDataSourceProvider(plan_node); }LakeDataSource一次湖表扫描的完整生命周期LakeDataSource是实际执行单次扫描的数据源其生命周期方法完整覆盖了打开→取数→复用→关闭open()初始化槽位描述符、解析列访问路径ColumnAccessPath、执行常量谓词求值、重建扫描谓词、初始化湖表读取器get_next()从投影迭代器projection iterator拉取 Chunk并在其上执行非下推谓词树的过滤与表达式的求值close()/release_for_reuse()/reuse()负责统计更新、资源释放以及为 prepared-split 场景下的读取器复用提供支持。从源码结构lake_connector.cpp 第 160 行起可以看出扫描参数的准备被拆成一系列职责单一的内部函数构成了清晰的初始化流水线open() └─ rebuild_scan_conjuncts() # 基于 ScanConjunctsManager 重建扫描谓词 └─ build_scan_range() # 构造扫描范围 └─ init_tablet_reader() # 初始化湖表读取器 ├─ get_tablet() # 解析 tablet_id/version获取版本化 Tablet ├─ init_global_dicts() # 建立槽位到存储列的全局字典映射 ├─ init_unused_output_columns() ├─ init_reader_params() # 汇总读取器参数 ├─ init_scanner_columns() # 确定扫描列与读取列 ├─ init_column_access_paths() / prune_schema_by_access_paths() └─ new_reader() projection iteratorget_tablet()中有两套 schema 获取路径值得注意若 FE 已启用快速 schema 演进 v2TLakeScanNode携带schema_key则通过table_schema_service()-get_schema_for_scan()获取 schema否则回退到旧的_tablet.get_schema()路径这一设计保证了滚动升级期间 FE/BE 版本的兼容性。LakeDataSourceProvider扫描范围到 morsel 队列的转换LakeDataSourceProvider负责把 FE 下发的TScanRangeParams转换为执行引擎可消费的 morsel 队列lake_connector.cpp 第 1843 行起。其中包含 BE 侧的动态分区裁剪init()时准备好的分区谓词上下文_partition_conjunct_ctxs会在convert_scan_range_to_morsel_queue_builder()中通过prune_scan_ranges_by_partition_conjuncts()提前裁剪掉不满足分区谓词的扫描范围从而减少无效 IO。该 Provider 还暴露了一组与湖表物理拆分相关的能力位could_split()/could_split_physically()是否可拆分、是否可物理拆分enable_lake_prepared_physical_split_scan()是否启用 prepared 物理拆分扫描读取器复用优化的开关sorted_by_keys_per_tablet()/output_chunk_by_bucket()/is_asc_hint()来自TLakeScanNode的排序与输出顺序提示用于下游算子优化。关键机制一谓词下推与运行时过滤器处理湖表扫描的性能很大程度取决于谓词能下推多深。LakeDataSource用ScanConjunctsManager统一管理扫描谓词rebuild_scan_conjuncts()构造ScanConjunctsManagerOptions传入列名、运行时过滤器、max_scan_key_num等随后调用parse_conjuncts()完成谓词分类与解析init_reader_params()从谓词树中取出可下推部分pushdown_pred_root交给存储层不可下推部分non_pushdown_pred_root保留在数据源上层执行谓词树通过PredicateTree统一组织OlapPredicateParser::can_pushdown()决定每个谓词节点的归属。一个值得注意的设计是parse_runtime_filters()被直接跳过返回Status::OK()// parse_runtime_filters is used to generate min-max predicates from runtime filters, while LakeDataSource already // generates predicates by ScanConjunctsManager, so skip parse_runtime_filters to make the parse logic is consistent // to the share-nothing mode. Status parse_runtime_filters(RuntimeState* state) override { return Status::OK(); }注释解释得很清楚由于LakeDataSource已经通过ScanConjunctsManager生成了谓词跳过parse_runtime_filters可以避免重复生成 min-max 谓词使解析逻辑与共享无状态share-nothing模式保持一致。对于运行时过滤器模块还实现了迟到运行时过滤器重初始化late runtime filter reinit机制。capture_runtime_filter_snapshots()记录每个过滤器的描述符、是否流式构建、是否到达及版本号needs_late_runtime_filter_reinit()对比新旧快照判断是否需要在读取器打开后重新初始化若需要则走reinit_reader_with_late_runtime_filters()路径——先禁用依赖运行时过滤器的 prepared 缓存重置读取器重建谓词并重新打开 tablet reader从而让新到达的运行时过滤器也能参与下推。关键机制二prepared 物理拆分与读取器复用为了降低湖表扫描的重复初始化成本ConnectorLake 实现了 prepared-split 机制扫描在初次打开时基于行集rowset/段segment元数据做种子准备seed prepare生成可复用的prepared_tablet_read_state与prepared_segment_read_state后续子拆分child split直接复用这些状态。相关代码路径集中在 lake_connector.cpp 的匿名命名空间与apply_child_split_context()中is_initial_coarse_split()/is_pre_refinement_coarse_split()根据LakeSplitContext::RowidRangeSource区分拆分来源初始粗拆分、预细化粗拆分、已细化拆分trim_coarse_split_by_pruned_range()用段级的已裁剪扫描范围pruned_scan_range对粗拆分的 rowid 范围做二次裁剪can_reuse_prepared_segment_for_child_split()判断子拆分能否直接复用 prepared 段读取状态reuse()/reopen_reader()LakeDataSource可通过ReusableReaderKey内部持有prepared_tablet_read_state判断当前 morsel 是否与可复用读取器匹配匹配时直接reopen_reader()而无需完整重建读取器。与之配套的还有一组细粒度的 profile 计数器如_lake_prepared_scan_rows_counter、_lake_reusable_segment_iter_reused_counter、_lake_late_rf_reinit_counter、_lake_seed_*系列计时器用于观测 prepared 路径的命中率与耗时拆解便于在真实负载下验证优化收益。关键机制三湖表延迟物化上下文ConnectorLake 的另一半职责——湖表延迟物化状态——由 lake_global_late_materialization_context.h 中的LakeScanLazyMaterializationContext承担。它继承自compute_env层的GlobalLateMaterilizationContext以 plan node 为粒度在查询运行态中按需创建get_or_create_ctx核心能力包括capture_rowsets()在读取器初始化时捕获当前扫描涉及的 rowsets、版本号以及缓存选项LakeScanCacheOptionsuse_page_cache、fill_data_cache、fill_metadata_cache、skip_disk_cacheget_rowset()根据动态 rowset iddrssid反查 rowset 与段索引供延迟物化阶段按需拉取对应数据内部以shared_mutex保护_rowsets、_versions、_cache_options三个映射保证多线程扫描场景下的并发安全。启用条件在 lake_connector.cpp 的init_tablet_reader()中当TLakeScanNode.enable_global_late_materialization为真时从query_runtime_state()-global_late_materialization_ctx_mgr()获取或创建上下文并设置扫描节点信息与缓存选项。延迟物化的核心收益是只有真正被查询用到的列才从湖存储远端对象存储物化到内存避免为不参与后续算子的列付出不必要的 IO 与解压开销。向量检索与 CACHE SELECT 等扩展能力从LakeDataSource的成员与open()逻辑可以推断湖表扫描还承载了若干高级特性向量索引检索open()中解析vector_search_options启用 ANN近似最近邻搜索支持refine_distance其废弃前身为use_ivfpq为兼容滚动升级仍会被识别、vector_limit_k、query_vector、vector_range、result_order等参数解析查询向量元素时使用不抛异常的StringParser对1e308这类超出 float 范围的输入返回InvalidArgument而非让 BE 崩溃CACHE SELECT 与主键索引预热当lake_cache_select_in_physical_way开启且查询选中了全部主键列时has_all_pk_columns_selected()会注册sst_warmup_fn回调在读取时异步预热云原生主键表的持久化索引 SST 文件warmup_pk_index_sst_files()JSON 列访问路径_inherit_default_value_from_json()支持从 JSON 父列的默认值中按访问路径如profile.level→$.level提取子字段默认值服务于 schema 按访问路径裁剪后的列默认值补全。测试与验证ConnectorLake 的单元测试位于 be/test/connector/lake 目录与模块边界清单中的allowed_test_targets: [connector_lake_test]对应lake_data_source_test.cpp覆盖LakeDataSource的打开、取数、复用等行为lake_global_late_materialization_context_test.cpp覆盖延迟物化上下文的 rowsets 捕获与按 drssid 反查逻辑lake_test_main.cpp测试入口。此外LakeDataSource特意保留了TEST_tablet_schema()与TEST_params()两个测试专用访问器并在LakeDataSourceProvider中提供set_lake_tablet_manager()注入点说明模块在设计之初就为单测隔离做好了准备。对模块边界本身的校验则可运行python3 build-support/check_be_module_boundaries.py --mode full。小结与阅读路径ConnectorLake 是 StarRocks 共享数据架构在 BE 侧的落地点之一LakeConnector通过Connector::LAKE类型被bootstrap_builtin_connectors()注册进连接器注册表LakeDataSource/LakeDataSourceProvider完成扫描执行、谓词下推、prepared 拆分复用与运行时过滤器重初始化LakeScanLazyMaterializationContext支撑湖表列的按需物化。而 be/module_boundary_manifest.json 中的connectorlake条目则以可机器校验的形式锁定了该模块的依赖边界。建议的后续阅读顺序be/src/connector/lake/lake_connector.h先建立类骨架认知be/src/connector/lake/lake_connector.cpp按open → init_tablet_reader → get_next → close的生命周期顺序精读be/src/connector/lake/lake_global_late_materialization_context.h理解延迟物化状态管理be/module_boundary_manifest.jsonconnectorlake条目理解模块边界规则be/src/module/connector_bootstrap.cpp理解内置连接器的注册方式be/test/connector/lake通过单测观察行为期望。【免费下载链接】starrocksThe worlds 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),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。