
SeaTunnel StarRocks Source 连接ector 实战查询计划获取、Tablet 切分与 BE 直读原理详解【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文以 Apache SeaTunnel 的 StarRocks 源端连接器StarRocks Source为主题完整覆盖其全部配置项、单表/多表读取示例与 Tablet 级切分策略并结合仓库源码剖析其底层实现链路从向 FE 请求查询计划Query Plan到按request_tablet_size将 Tablet 切分为 Split再到 Reader 通过 Thrift 直连 BE 批量拉取数据的全过程帮助你既会用这个连接器也知其所以然。一、概述与核心能力StarRocks Source 连接器用于从 StarRocks 读取外部数据源的数据。官方文档对其内部实现的描述是先从前端FE获取查询计划再将查询计划作为参数下发给 BE 节点最后从 BE 节点获取数据结果见 StarRocks.md。从源码可以印证这条链路StarRocksSource 实现SeaTunnelSource接口其getBoundedness()返回BOUNDED因此该连接器只支持批模式batch不支持流模式这与文档 Key features 中的能力清单一致StarRocksQueryPlanReadClient 负责向 FE 请求查询计划并按request_tablet_size生成QueryPartitionStarRocksBeReadClient 封装 Thrift 客户端直连 BE 执行open_scanner/get_next读取 Arrow 数据。当前支持的 Connector V2 能力如下与文档一致能力支持情况批模式batch支持流模式stream不支持exactly-once不支持schema projection字段投影支持并行度parallelism支持自定义 split支持二、配置项Options完整说明以下参数表完整继承自官方文档并结合 StarRocksSourceOptions 与 StarRocksBaseOptions 中Options.key(...)的声明核对过默认值与描述nametyperequireddefault valuedescriptionnodeUrlslistyes-StarRocks FE HTTP 地址格式[fe_ip:fe_http_port, ...]usernamestringyes-StarRocks 用户名passwordstringyes-StarRocks 密码databasestringyes-StarRocks 数据库名tablestringno-单表名未配置table_list时必填table_listarrayno-多表读取列表未配置table时必填。每个条目可单独定义schema与scan_filterschemaconfigno-输出 schema单表模式配置在顶层多表模式配置在每个table_list条目内scan_filterstringno透传给 StarRocks 的源端过滤表达式request_tablet_sizeintnoInteger.MAX_VALUE单个 SeaTunnel split 内允许的最大 Tablet 数值越小 split 越多scan_connect_timeout_msintno1000连接 StarRocks BE 进行 scan 的超时时间毫秒scan_query_timeout_secintno3600查询超时时间秒-1表示不超时scan_keep_alive_minintno10查询任务的保活时间分钟scan_batch_rowsintno1024每次从 BE 读取的最大行数scan_mem_limitlongno1073741824单个 BE 查询允许使用的最大内存字节默认 1GBmax_retriesintno3向 StarRocks 发起请求的重试次数scan.params.*stringno-额外的 BE scan 参数发送前会去掉scan.params.前缀其中必填项的约束由 StarRocksSourceFactory 的optionRule()以 OptionRule 方式声明nodeUrls、username、password、database为 requiredtable与table_list为 exclusive二选一不可同时配置。nodeUrls [list]StarRocks 集群的 FE 地址格式为[fe_ip:fe_http_port, ...]。注意这里是FE 的 HTTP 端口常见为 8030用于获取查询计划实际数据读取走的是 BE 的 Thrift 端口由查询计划返回的 routing 信息决定无需额外配置。从源码结构看该配置天然支持高可用StarRocksQueryPlanReadClient在发起请求前会先对nodeUrls做Collections.shuffle()打乱顺序然后依次尝试各个 FE 节点注释明确写着 shuffle nodeUrls to ensure support for both random selection and high availability单个节点失败会自动切换下一个全部失败才抛出QUEST_QUERY_PLAN_FAILED异常。username / password / databaseStarRocks 用户名、密码与数据库名。用户名密码会被用于两个环节向 FE 请求查询计划时以Authorization: Basic base64(user:pass)头携带向 BE 打开 scanner 时写入TScanOpenParams的user/passwd字段见StarRocksBeReadClient.openScanner。table [string] 与 table_list [array]单表与多表模式table与table_list二选一使用table时schema配置在 source 顶层使用table_list时schema配置在每个表条目内部。配置解析由 StarRocksSourceTableConfig 完成of(config)先检查是否存在CatalogOptions.TABLE_LIST存在则把列表中每个Map转成ReadonlyConfig后逐个解析为StarRocksSourceTableConfig各自持有 table 名、CatalogTable和scanFilter否则解析顶层配置生成单元素列表。因此多表模式只是把单表模式统一收敛为表配置列表后续切分、读取逻辑完全一致。每个表条目支持table、schema、scan_filter三个字段可以实现不同表使用不同字段投影和不同过滤条件。schema [config]定义 SeaTunnel 要输出的 StarRocks 行的 schema。字段投影schema projection在生成查询 SQL 时生效StarRocksQueryPlanReadClient.genQuerySql()会取CatalogTable中SeaTunnelRowType的字段名拼接成select col1,col2,... from \db.table字段为空时才退化为select *。也就是说schema 中声明的字段直接决定了查询计划与最终输出的列集合。schema 的详细语法请参考 SeaTunnel 的 Schema 特性文档。示例schema { fields { name string age int } }scan_filter [string]查询的过滤表达式会被透明透传给 StarRocks由 StarRocks 完成源端数据过滤。例如tinyint_1 100实现上scan_filter直接拼进查询 SQL 的where子句genQuerySql中filter scanFilter.isEmpty() ? : where scanFilter随查询计划一起下发StarRocks 在 BE 侧执行过滤SeaTunnel 不额外承担过滤计算。request_tablet_size [int]单个 split 包含的 Tablet 数上限。该值越小生成的 split 越多引擎侧并行度越高但对 StarRocks 的访问压力也越大默认值为Integer.MAX_VALUE即不做切分限制。文档中给出的 Tablet 分布示例很好地说明了切分逻辑。假设集群中 Tablet 分布如下be_node_1 tablet[1, 2, 3, 4, 5] be_node_2 tablet[6, 7, 8, 9, 10] be_node_3 tablet[11, 12, 13, 14, 15]若未设置request_tablet_size无限制split 按 BE 自然生成partition[0] read data of tablet[1, 2, 3, 4, 5] from be_node_1 partition[1] read data of tablet[6, 7, 8, 9, 10] from be_node_2 partition[2] read data of tablet[11, 12, 13, 14, 15] from be_node_3若设置request_tablet_size3每个 split 最多 3 个 Tabletpartition[0] read data of tablet[1, 2, 3] from be_node_1 partition[1] read data of tablet[4, 5] from be_node_1 partition[2] read data of tablet[6, 7, 8] from be_node_2 partition[3] read data of tablet[9, 10] from be_node_2 partition[4] read data of tablet[11, 12, 13] from be_node_3 partition[5] read data of tablet[14, 15] from be_node_3这与源码实现完全对应StarRocksQueryPlanReadClient.tabletsMapToPartition()先按 BE 聚合 Tablet 集合去重后再以requestTabletSize为步长切片每个切片生成一个QueryPartition携带 database、table、BE 地址、tablet 集合与序列化后的查询计划。值得注意的是切分严格以 BE 为边界——同一个 split 的 Tablet 一定属于同一个 BE因为StarRocksBeReadClient是按 BE 地址建立 Thrift 连接打开 scanner 的跨 BE 的 Tablet 无法在同一个 scanner 中读取。此外从源码结构看selectBeForTablet()在 Tablet 有多个可路由 BE副本时会优先选择当前负载已分配 Tablet 数最少的 BE起到一定的副本负载均衡作用。scan_connect_timeout_ms [int]与 BE 建立 Thrift 连接的超时时间毫秒默认 1000。实现上该值同时作为TSocket的连接超时与 socket 超时new TSocket(ip, port, connectTimeoutMs, connectTimeoutMs)因此网络较差的跨机房场景建议适当调大。scan_query_timeout_sec [int]单个查询的超时时间秒默认 36001 小时-1表示不限制。该值通过TScanOpenParams.setQuery_timeout()传递给 BE。scan_keep_alive_min [int]查询任务在 BE 侧的保活时长分钟默认 10文档建议设置为不小于 5 的值。实现中该值会被截断到Short.MAX_VALUE上限后写入TScanOpenParams.setKeep_alive_min()。保活时间的意义是当某个 split 因为执行慢而长时间没有拉取数据时BE 侧的 scanner 上下文不会过早回收。scan_batch_rows [int]单次从 BE 拉取的最大行数默认 1024。增大该值可以减少引擎与 StarRocks 之间建立的连接/往返次数从而缓解网络延迟带来的开销但该值过大会增加单批内存占用需与scan_mem_limit配合评估。scan_mem_limit [long]单个查询在 BE 节点上允许使用的最大内存空间字节默认 10737418241GB通过DEFAULT_SCAN_MEM_LIMIT 1024 * 1024 * 1024L定义于StarRocksSourceOptions并写入TScanOpenParams.setMem_limit()。max_retries [int]向 StarRocks 请求这里是向 FE 请求查询计划的失败重试次数默认 3。StarRocksQueryPlanReadClient使用RetryUtils.RetryMaterial(maxRetries, true, exception - true, 1000ms)构建重试材料即任何异常都重试、重试间隔固定 1000ms。scan.params.* [string]透传给 BE scan 过程的额外参数。SourceConfig的构造函数会遍历配置 Map收集所有以scan.params.为前缀的键值对去掉前缀并转成小写后存入sourceOptionProps最终通过TScanOpenParams.setProperties()整体下发给 BE。例如配置scan.params.scanner_thread_pool_thread_num 3会下发给 BE 的属性键为scanner_thread_pool_thread_num。三、完整作业示例示例 1单表读取source { StarRocks { nodeUrls [starrocks_e2e:8030] username root password database test table e2e_table_source scan_batch_rows 10 max_retries 3 schema { fields { BIGINT_COL BIGINT LARGEINT_COL STRING SMALLINT_COL SMALLINT TINYINT_COL TINYINT BOOLEAN_COL BOOLEAN DECIMAL_COL DECIMAL(20, 1) DOUBLE_COL DOUBLE FLOAT_COL FLOAT INT_COL INT CHAR_COL STRING VARCHAR_11_COL STRING STRING_COL STRING DATETIME_COL TIMESTAMP DATE_COL DATE } } scan.params.scanner_thread_pool_thread_num 3 } }该示例覆盖了全部常用配置FE HTTP 地址starrocks_e2e:8030、认证信息、单表 完整 schema 字段声明、批次行数与重试次数以及一条scan.params.*透传参数。示例 2多表读取table_listsource { StarRocks { nodeUrls [starrocks_e2e:8030] username root password database test table_list [ { table e2e_table_source schema { fields { BIGINT_COL BIGINT LARGEINT_COL STRING SMALLINT_COL SMALLINT TINYINT_COL TINYINT BOOLEAN_COL BOOLEAN DECIMAL_COL DECIMAL(20, 1) DOUBLE_COL DOUBLE FLOAT_COL FLOAT INT_COL INT CHAR_COL STRING VARCHAR_11_COL STRING STRING_COL STRING DATETIME_COL TIMESTAMP DATE_COL DATE } } }, { table e2e_table_source_2 schema { fields { BIGINT_COL_2 BIGINT LARGEINT_COL_2 STRING SMALLINT_COL_2 SMALLINT TINYINT_COL_2 TINYINT BOOLEAN_COL_2 BOOLEAN DECIMAL_COL_2 DECIMAL(20, 1) DOUBLE_COL_2 DOUBLE FLOAT_COL_2 FLOAT INT_COL_2 INT CHAR_COL_2 STRING VARCHAR_11_COL_2 STRING STRING_COL_2 STRING DATETIME_COL_2 TIMESTAMP DATE_COL_2 DATE } } }] scan_batch_rows 10 max_retries 3 scan.params.scanner_thread_pool_thread_num 3 } }多表模式下每张表各自声明schema还可以各自声明scan_filterscan_batch_rows、max_retries、scan.params.*等作业级参数仍然配置在顶层、对所有表生效。四、源码级运行流程剖析结合仓库源码整个 StarRocks Source 的作业运行可以分为三个阶段。阶段一Split 枚举——从 FE 获取查询计划并切分入口是 StartRocksSourceSplitEnumerator。其run()流程为从pendingTables由table/table_list解析出的表列表中逐张取表调用getStarRocksSourceSplit(table)生成 split 列表对每张表StarRocksQueryPlanReadClient.findPartitions(table)完成三件事genQuerySql(table)按 schema 字段与scan_filter拼装select ... from \db.table [where ...]getQueryPlan(...)向http://{fe}/api/{db}/{table}/_query_plan发起 POSTBasic 认证请求体为{sql: ...}FE 返回包含 opaquedQueryPlan 与 Tablet→BE routing 信息的QueryPlanJSONselectBeForTablet()tabletsMapToPartition()按上文所述的负载均衡与request_tablet_size切分规则生成QueryPartition列表每个QueryPartition被包装为一个 StarRocksSourceSplitsplitId 取 partition 的 hashCode 字符串所有 split 生成后addPendingSplit()按 splitId 排序以全局轮询方式assignCount % readerCount预分配到各个 reader再由assignSplit()下发最后向所有 reader 发送signalNoMoreSplits。assignCount、pendingTables、pendingSplit三者会被snapshotState()快照进StarRocksSourceState保证 failover 恢复后分配语义连续addSplitsBack()则支持 reader 失败时把 split 回传给指定 subtask 重新分配。阶段二Reader 读取——Thrift 直连 BE 批量拉取StarRocksSourceReader 维护一个按 BE 地址ip:port索引的StarRocksBeReadClient连接池clientsPools对同一 BE 的多个 split 复用同一个客户端实例。每次pollNext取出一个 pending split 后获取/创建对应 BE 的客户端首次创建时以scan_connect_timeout_ms建立TSocketclient.openScanner(partition, rowType)构造TScanOpenParams写入 tablet 集合、opaqued 查询计划、database/table/user/passwd、batch_sizescan_batch_rows、keep_alive_min、query_timeout、mem_limit以及scan.params.*透传的 properties调用open_scanner得到context_id循环client.hasNext()/getNext()以context_id offset调用get_next拉取 Arrow 格式批次TScanBatchResult经由ArrowToSeatunnelRowReader转换为SeaTunnelRow后输出全部 split 读完且收到NoMoreSplits事件后调用context.signalNoMoreElement()通知引擎该 source 结束批模式的正常终止方式。StarRocksSource的getProducedCatalogTables()会返回所有表的CatalogTable使下游能够感知每个表输出的精确 schema 类型这也是schema声明必须完整的原因。阶段三重试与异常边界FE 请求查询计划失败时按max_retries重试重试耗尽或响应为空时抛出StarRocksConnectorException(QUEST_QUERY_PLAN_FAILED)单个 FE 节点失败会记录错误日志并继续尝试下一个 FE 节点shuffle 后的顺序因此配置多个 FE 节点即可获得查询计划阶段的高可用BE 侧open_scanner/get_next返回非 OK 状态码时抛出SCAN_BE_DATA_FAILED异常由引擎层按作业失败策略处理。相关错误码定义在 StarRocksConnectorErrorCode。五、参数调优与适用场景建议基于上述实现给出几条可落地的调优方向均为机制推断具体取值需结合实际集群压测并行度与 StarRocks 压力的权衡默认request_tablet_sizeInteger.MAX_VALUE时 split 数等于 BE 数每个 BE 一个 split如果引擎并行度高于 BE 数调小request_tablet_size可以增加 split 数、提升并行利用率但同时会打开更多 BE scanner注意scan_mem_limit与 BE 内存的匹配网络延迟敏感场景调大scan_batch_rows如 4096/8192可以减少get_next往返次数跨机房部署时同步调大scan_connect_timeout_ms避免TSocket建连超时慢节点保护scan_keep_alive_min保证 split 长时间不拉取时 BE 侧 scanner 上下文不失效建议保持默认 10 分钟以上scan_query_timeout_sec在超大表全量读取场景可设为-1避免误超时BE 侧参数透传scan.params.*是官方预留的 BE 调优通道示例中的scan.params.scanner_thread_pool_thread_num即用于控制 BE 端 scan 线程池线程数其余 BE 支持的 scan 属性同理透传字段裁剪与源端过滤优先使用schema声明所需字段直接缩减查询计划列集合与scan_filter下推过滤条件由 StarRocks 在 BE 侧执行减少网络与序列化开销多表场景下各表可独立配置这两项。六、小结StarRocks Source 是 SeaTunnel 中一个典型的查询计划 直读存储型批处理连接器FE 负责 SQL 语义解析与 Tablet 路由SeaTunnel 负责 Tablet 级切分与并行调度BE 负责最终的数据扫描。理解request_tablet_size的切分语义、scan_batch_rows/scan_mem_limit的资源边界、以及scan.params.*的透传机制是把这个连接器用好、调优到位的关键。如需查看该连接器的完整变更历史可参考 connector-starrocks changelog例如 Support multi starrocks source、Support StarRocks Fe Node HA 等条目对应本文源码中的多表模式与 FE 高可用逻辑。更多连接器背景可参阅 连接ector 概览 与 Schema 特性。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考