资讯动态

SeaTunnel StarRocks Sink 连接器详解:Stream Load 批量写入、自动建表与 CDC 同步实战

发布时间:2026/9/20 1:34:47 来源:尧图企业网站定制
SeaTunnel StarRocks Sink 连接器详解Stream Load 批量写入、自动建表与 CDC 同步实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 StarRocks 数据接收器Sink用于将数据批量写入 StarRocks内部通过缓存与 Stream Load 通道实现批导入支持批BATCH与流STREAMING两种模式。本文以 docs/zh/connectors/sink/StarRocks.md 为骨架结合 connector-starrocks 模块源码完整讲解全部接收器选项、自动建表模板、SaveMode 语义、CDC 变更事件同步与 Zeta 定时刷新机制读完后你可以直接在 SeaTunnel 中配置出一个生产可用的 StarRocks Sink 作业。引擎支持与特性矩阵StarRocks Sink 在以下引擎上可用SparkFlinkSeaTunnel Zeta在 Connector V2 特性矩阵 中该连接器的能力标记如下精确一次Exactly-once不支持未勾选CDC 变更事件同步支持多表写入Multi-table支持定时刷新支持仅 Zeta 引擎见下文需要特别注意的是尽管 Stream Load 本身基于 label 具备幂等能力但当前官方文档明确将精确一次排除在已支持能力之外StarRocks Sink 实际为至少一次At-least-once语义任务重启后可能重复提交数据。工作原理缓存 Stream Load 批导入该接收器用于将数据写入 StarRocks内部实现采用缓存机制当满足刷新条件后通过 HTTP 协议调用 StarRocks 的 Stream Load 接口完成批量导入。从源码看完整写入链路为StarRocksSinkWriter.write()将每条SeaTunnelRow序列化为字符串CSV 或 JSON交给StarRocksSinkManager.write()缓存见 StarRocksSinkWriter.javaStarRocksSinkManager维护batchList当缓存行数达到batch_max_rows、字节数达到batch_max_bytes或触发 checkpoint/timer flush 时执行flush()见 StarRocksSinkManager.javaStarRocksStreamLoadVisitor.doStreamLoad()向http://nodeUrl/api/{database}/{table}/_stream_load发起 PUT 请求携带 label、Basic Auth 与 data_desc 参数见 StarRocksStreamLoadVisitor.java。其中getAvailableHost()会按轮询方式遍历nodeUrls先探测连通性再选择可用 FE 节点发起导入StarRocksStreamLoadVisitor.java#L186-L196。请求头中会强制加入strip_outer_arraytrue、Expect: 100-continue、Content-Type: application/x-www-form-urlencoded并将starrocks.config中的键值逐一注入 HTTP 头StarRocksStreamLoadVisitor.java#L367-L397。失败重试与 label 幂等flush()中最多重试max_retries次每次重试前按min(retry_backoff_multiplier_ms * i, max_retry_backoff_ms)退避等待当 Stream Load 返回Label Already Exists或Fail且消息指明 label 已被占用时会先通过/api/{db}/get_load_state?label...查询 label 的最终状态VISIBLE/COMMITTED视为成功避免在结果未知时重复提交数据仅当 StarRocks 明确报告 label 处于ABORTED状态时才更换新 label 重新提交StarRocksSinkManager.java#L106-L176。依赖准备对于 Spark/Flink你需要下载 MySQL JDBC 驱动 jar即mysql-connector-java可从 Maven 中央仓库获取并放置到目录${SEATUNNEL_HOME}/plugins/下。对于 SeaTunnel Zeta同样需要 MySQL JDBC 驱动 jar但放置目录为${SEATUNNEL_HOME}/lib/。该驱动用于通过base-url建立 JDBC 连接执行建表DDL、查询 schema、SaveMode 预处理等操作StarRocksSinkWriter.applySchemaChange()中会显式加载com.mysql.cj.jdbc.Driver并建立 JDBC 连接执行 schema 变更见 StarRocksSinkWriter.java#L104-L119。接收器选项下表汇总了 StarRocks Sink 的全部配置项与 StarRocksSinkOptions.java 中的 Option 定义一一对应名称类型是否必须默认值说明nodeUrlslist是-StarRocks 集群地址格式为[fe_ip:fe_http_port, ...]用于 Stream Load 数据写入base-urlstring是-JDBC URL 样式的连接信息。如jdbc:mysql://localhost:9030/、jdbc:mysql://localhost:9030或jdbc:mysql://localhost:9030/db用于建表与 schema 相关 DDLusernamestring是-目标 StarRocks 用户名passwordstring是-目标 StarRocks 密码databasestring是-目标 StarRocks 表所在的数据库名称tablestring否-目标 StarRocks 表名。如果没有设置则表名与上游表名相同labelPrefixstring否-StarRocks Stream Load 作业标签前缀batch_max_rowslong否1024批量写入时当缓存行数达到batch_max_rows、字节数达到batch_max_bytes或时间达到checkpoint.interval时数据会刷新到 StarRocksbatch_max_bytesint否5 * 1024 * 1024批量写入时当缓存行数达到batch_max_rows、字节数达到batch_max_bytes或时间达到checkpoint.interval时数据会刷新到 StarRocksmax_retriesint否-数据写入 StarRocks 失败后的重试次数retry_backoff_multiplier_msint否-用作生成下一次退避延迟的乘数max_retry_backoff_msint否-向 StarRocks 发送重试请求前的最大等待时长enable_upsert_deleteboolean否false是否开启 upsert/delete 事件同步仅支持主键模型Primary Key表save_mode_create_templatestring否见下方默认模板自动建表模板starrocks.configmap否-Stream Loaddata_desc参数http_socket_timeout_msint否180000HTTP socket 超时时间默认为 3 分钟schema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST同步任务启动前针对目标端已存在的表结构选择不同处理方式data_save_modeEnum否APPEND_DATA同步任务启动前针对目标端已存在的数据选择不同处理方式custom_sqlString否-当data_save_mode设置为CUSTOM_PROCESSING时必须配置。该 SQL 会在同步任务启动前执行源码层面的参数约束见 StarRocksSinkFactory.java#L54-L84必填项username、password、database、base-url、nodeUrlscustom_sql为条件参数仅当data_save_mode CUSTOM_PROCESSING时要求配置否则视为非法组合table_options通过StarRocksTableOptionsConditionExtension做扩展条件校验这些规则在--check配置校验与作业提交阶段就会触发而非等到 StarRocks 执行 DDL 时才失败。save_mode_create_templateStarRocks 数据接收器使用模板在需要的时候也可以修改模板并结合上游数据类型和结构生成表的创建语句来自动创建 StarRocks 表。当前仅在多表模式下有效。默认模板与源码 StarRocksSinkOptions.java#L45-L70 中定义的默认值一致CREATE TABLE IF NOT EXISTS ${database}.${table_name} ( ${rowtype_primary_key}, ${rowtype_fields} ) ENGINEOLAP PRIMARY KEY (${rowtype_primary_key}) COMMENT ${comment} DISTRIBUTED BY HASH (${rowtype_primary_key})PROPERTIES ( replication_num 1 )在模板中添加自定义字段比如加上id字段的修改模板如下CREATE TABLE IF NOT EXISTS ${database}.${table_name} ( id, ${rowtype_fields} ) ENGINE OLAP COMMENT ${comment} DISTRIBUTED BY HASH (${rowtype_primary_key}) PROPERTIES ( replication_num 1 );StarRocks 数据接收器根据上游数据自动获取相应的信息来填充模板并且会移除rowtype_fields中的 id 字段信息。使用此方法可用来为自定义字段修改类型及相关属性。可以使用的占位符有database上游数据模式的库名称table_name上游数据模式的表名称rowtype_fields上游数据模式的所有字段信息连接器会将字段信息自动映射到 StarRocks 对应的类型rowtype_primary_key上游数据模式的主键信息结果可能是列表rowtype_unique_key上游数据模式的唯一键信息结果可能是列表comment上游数据模式的注释信息从源码看模板中字段的生成由StarRocksSaveModeUtil.columnToConnectorType()完成会输出列名 类型 NULL/NOT NULL [COMMENT ...]四段式列定义见 StarRocksSaveModeUtil.java#L106-L127。table [string]使用选项参数database和table-name自动生成 SQL并接收上游输入数据写入 StarRocks 中。此选项与query是互斥的且具有更高的优先级。table选项参数可以填入任意表名这个名字最终会被用作目标表的表名并且支持变量${table_name}${schema_name}。替换规则如下${schema_name}将替换传递给目标端的 SCHEMA 名称${table_name}将替换传递给目标端的表名。例如test_${schema_name}_${table_name}_testsink_sinktabless_${table_name}补充实现细节当table未配置时StarRocksSinkFactory.createSink()会自动使用上游CatalogTable的表名填充并基于sinkConfig.getDatabase()与解析后的表名重写目标TableIdentifier见 StarRocksSinkFactory.java#L92-L116。schema_save_mode [Enum]在同步任务打开之前针对目标端已存在的表结构选择不同的处理方法。可选值有RECREATE_SCHEMA不存在的表会直接创建已存在的表会删除并根据参数重新创建CREATE_SCHEMA_WHEN_NOT_EXIST忽略已存在的表不存在的表会直接创建默认值ERROR_WHEN_SCHEMA_NOT_EXIST当有不存在的表时会直接报错IGNORE忽略对表的处理data_save_mode [Enum]在同步任务打开之前针对目标端已存在的数据选择不同的处理方法。可选值有DROP_DATA保存数据库结构但是会删除表中存量数据APPEND_DATA保存数据库结构和相关的表存量数据默认值CUSTOM_PROCESSING自定义处理ERROR_WHEN_DATA_EXISTS当对应表存在数据时直接报错custom_sql [String]当data_save_mode设置为CUSTOM_PROCESSING时必须同时设置custom_sql参数。custom_sql的值为可执行的 SQL 语句在同步任务开启前 SQL 将会被执行。table_options [Map]Sink 在 SaveMode 自动建表DDL时附加的表级属性。仅在schema_save_mode触发建表时生效例如CREATE_SCHEMA_WHEN_NOT_EXIST、RECREATE_SCHEMA不影响Stream Load 写入也不会对已存在表执行ALTER TABLE。在默认save_mode_create_template未配置或与内置默认值相同下table_options会合并进模板的PROPERTIES子句同名 key 以table_options为准。属性名请参考 StarRocks 官方CREATE TABLE文档中的 PROPERTIES 部分SeaTunnel 不做白名单非法属性由 StarRocks 执行 DDL 时报错。若配置了与内置默认值不同的save_mode_create_template则不能与table_options同时使用任务提交时校验失败此时请将属性直接写入模板。非法组合会在StarRocksSinkFactory的 option 规则阶段提前校验--check与作业提交而非仅在 StarRocks 执行 CREATE TABLE 时失败。实现佐证StarRocksSaveModeUtil.validateTableOptions()会检查table_options是否与自定义模板互斥以及 key/value 是否为空applyTableOptionsToCreateTableSql()通过SqlTableClauseMerger.merge()以双引号包裹的PROPERTIES格式将属性合并进建表 SQL见 StarRocksSaveModeUtil.java#L61-L94。示例sink { StarRocks { base-url jdbc:mysql://127.0.0.1:9030 nodeUrls [127.0.0.1:8030] username root password database test schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST table_options { replication_num 3 storage_format V2 } } }Zeta 定时刷新该引擎级能力仅由 Zeta 支持。可以在env块中配置sink.flush.interval使尚未达到batch_max_rows和batch_max_bytes的缓冲数据也能定时通过 StarRocks Stream Load 写出。Spark 和 Flink 不会触发该定时刷新。对应的回调注册位于StarRocksSinkWriter构造器中的context.registerFlushAction(this::timerFlush)见 StarRocksSinkWriter.java#L61-L73。注意StarRocks 定时刷新不提供基于 2PC 的精准一次语义StarRocks Sink 仍为至少一次语义任务重启后可能重复提交数据。如果业务场景适用可以使用具有确定性主键的 Primary Key 表吸收重复写入。env { job.mode STREAMING checkpoint.interval 300000 sink.flush.interval 5000 } sink { StarRocks { nodeUrls [starrocks-fe:8030] base-url jdbc:mysql://starrocks-fe:9030/mydb username root password database mydb table mytable batch_max_rows 10000 batch_max_bytes 104857600 } }数据类型映射StarRocks 数据类型与 SeaTunnel 数据类型的映射关系如下StarRocks 数据类型SeaTunnel 数据类型BOOLEANBOOLEANTINYINTTINYINTSMALLINTSMALLINTINTINTBIGINTBIGINTFLOATFLOATDOUBLEDOUBLEDECIMALDECIMALDATESTRINGTIMESTRINGDATETIMESTRINGSTRINGSTRINGARRAYSTRINGMAPSTRINGBYTESSTRING自动建表时的反向类型推断在自动建表场景schema_save_mode触发 DDL下连接器会按StarRocksSaveModeUtil.dataTypeToStarrocksType()将上游 SeaTunnel 类型转换为 StarRocks 建表类型见 StarRocksSaveModeUtil.java#L129-L176关键规则包括NULL、TIME→VARCHAR(8)STRING列长度大于 65533 或未指定长度时映射为STRING否则映射为VARCHAR(length)BYTES→STRINGDATE→DATETIMESTAMP→DATETIMEARRAYT→ARRAY元素类型元素类型递归转换DECIMAL(p, s)→Decimal(p, s)保留精度与小数位MAP、ROW→JSON支持导入的数据格式StarRocks 数据接收器支持的格式有 CSV 和 JSON 两种通过starrocks.config中的format指定枚举定义于 SinkConfig.java#L38-L41序列化器分别由 StarRocksCsvSerializer.java 与 StarRocksJsonSerializer.java 承担JSON默认格式。批次数据被拼接为[...]数组首尾加方括号、元素间以逗号分隔配合强制的strip_outer_arraytrue头参数导入CSV使用column_separator、row_delimiter支持\x01、\x02等转义由StarRocksDelimiterParser解析拼接行数据并在请求头中生成columns字段按列名顺序映射。任务示例简单示例该示例包含多种数据类型的数据写入且用户需要为目标端下游创建相应表env { parallelism 1 job.mode BATCH checkpoint.interval 10000 } source { FakeSource { row.num 10 map.size 10 array.size 10 bytes.length 10 string.length 10 schema { fields { c_map mapstring, arrayint c_array arrayint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_decimal decimal(16, 1) c_null null c_bytes bytes c_date date c_timestamp timestamp } } } } sink { StarRocks { nodeUrls [e2e_starRocksdb:8030] base-url jdbc:mysql://e2e_starRocksdb:9030/ username root password database test table e2e_table_sink batch_max_rows 10 starrocks.config { format JSON strip_outer_array true } } }支持写入 CDC 变更事件INSERT/UPDATE/DELETE示例sink { StarRocks { nodeUrls [e2e_starRocksdb:8030] base-url jdbc:mysql://e2e_starRocksdb:9030/ username root password database test table e2e_table_sink ... // 支持upsert/delete事件的同步需要将选项参数enable_upsert_delete设置为true仅支持表引擎为主键模型 enable_upsert_delete true } }CDC 事件如何传播StarRocksSinkOP.parse(RowKind)将INSERT、UPDATE_AFTER映射为UPSERT将DELETE、UPDATE_BEFORE映射为DELETE并写入固定列__op见 StarRocksSinkOP.java。开启enable_upsert_delete后序列化器会追加__op列CSV 模式下该列也会加入请求头columns声明中StarRocksStreamLoadVisitor.java#L372-L374。JSON 格式数据导入示例sink { StarRocks { nodeUrls [e2e_starRocksdb:8030] base-url jdbc:mysql://e2e_starRocksdb:9030/ username root password database test table e2e_table_sink batch_max_rows 10 starrocks.config { format JSON strip_outer_array true } } }CSV 格式数据导入示例sink { StarRocks { nodeUrls [e2e_starRocksdb:8030] base-url jdbc:mysql://e2e_starRocksdb:9030/ username root password database test table e2e_table_sink batch_max_rows 10 starrocks.config { format CSV column_separator \\x01 row_delimiter \\x02 } } }使用 save_mode 的示例sink { StarRocks { nodeUrls [e2e_starRocksdb:8030] base-url jdbc:mysql://e2e_starRocksdb:9030/ username root password database test table test_${schema_name}_${table_name} schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_modeAPPEND_DATA batch_max_rows 10 starrocks.config { format CSV column_separator \\x01 row_delimiter \\x02 } } }常见问题StarRocks Sink 支持自动建表吗支持。通过schema_save_mode参数控制建表行为CREATE_SCHEMA_WHEN_NOT_EXIST表不存在时创建已存在则跳过。RECREATE_SCHEMA每次任务启动时删除并重建表。ERROR_WHEN_SCHEMA_NOT_EXIST表不存在时抛出异常。IGNORE跳过所有建表逻辑。SeaTunnel 会根据上游 schema 自动推断 StarRocks 列类型推断规则详见上文自动建表时的反向类型推断小节。StarRocks Sink 是否支持 Upsert 和 DELETE 操作支持。设置enable_upsert_delete true可以传播 Upsert 和 DELETE 事件目标 StarRocks 表必须使用主键模型Primary Key。来自 CDC 数据源的 DELETE 事件在开启此选项后可正确传播。StarRocks Sink 中的 labelPrefix 是做什么的当前 StarRocks Sink 页面并未将精确一次列为已支持的 Connector 能力。labelPrefix用于控制 Sink 生成的 Stream Load label 前缀。从源码看每批次 label 的生成逻辑为若配置了labelPrefix则先拼接该前缀再追加一个UUID.randomUUID()形成{labelPrefix}{uuid}见 StarRocksSinkManager.java#L185-L191。保持此前缀稳定且全局唯一可以减少重试或任务重启时的 label 冲突sink { StarRocks { nodeUrls [starrocks-fe:8030] base-url jdbc:mysql://starrocks-fe:9030/ username root password database mydb table mytable labelPrefix unique-job-label } }正式契约请以本页的主要特性矩阵和labelPrefixoption 说明为准。StarRocks 列名是否区分大小写StarRocks 列名默认不区分大小写。请确认上游字段名与目标 StarRocks 列名的映射关系避免意外的字段不匹配。nodeUrls 和 base-url 有什么区别nodeUrlsStarRocks FE 节点的 HTTP 地址用于 Stream Load 数据写入。base-url指向 StarRocks FE 节点的 JDBC URL用于建表、查询 schema 等 DDL 操作。开启自动建表时两者均需配置。从架构上看nodeUrls走StarRocksStreamLoadVisitor的 HTTP PUT 通道base-url走 JDBC 通道DriverManager.getConnection两条通道职责分离。变更日志该连接器的历史变更记录维护在 docs/zh/connectors/changelog/connector-starrocks.md在文档站中通过ChangeLog /组件动态渲染升级版本或排查行为变化时建议同步查阅。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价