StarRocks Flink Connector 版本发布解读:1.2.8 至 1.2.15 兼容性矩阵、获取方式与关键能力演进

StarRocks Flink Connector 版本发布解读:1.2.8 至 1.2.15 兼容性矩阵、获取方式与关键能力演进 StarRocks Flink Connector 版本发布解读1.2.8 至 1.2.15 兼容性矩阵、获取方式与关键能力演进【免费下载链接】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 官方仓库中的 Flink Connector 发布记录 展开系统梳理 StarRocks Connector for Apache Flink以下简称 Flink connector的 JAR 命名规则与获取方式、各版本与 Flink/StarRocks/Java/Scala 的兼容性矩阵并逐版本解读 1.2.8 至 1.2.15 的新特性、改进与缺陷修复。读完本文你可以为现有 Flink 环境选定合适的 connector 版本并理解 exactly-once 事务、Merge Commit、多表事务 Stream Load、读端复杂类型支持等关键能力在不同版本中的落地过程与配套参数。一、获取 Flink Connector JAR 的方式与命名规则Flink connector 以独立仓库starrocks-connector-for-apache-flink发布产物发布到 Maven Central。根据发布记录获取 JAR 的方式有三种直接下载从 Maven Central 仓库com.starrocks组下下载已编译好的 JAR 文件Maven 依赖在项目pom.xml中添加依赖后由 Maven 拉取自行编译下载 connector 源码执行sh build.sh flink_version编译出 JAR 包。JAR 文件的命名格式与 Flink 版本强相关Flink 版本JAR 命名格式Flink 1.15 及以上flink-connector-starrocks-${connector_version}_flink-${flink_version}.jarFlink 1.15 之前flink-connector-starrocks-${connector_version}_flink-${flink_version}_${scala_version}.jar从 Flink 1.15 起Flink 不再强制要求 Scala 后缀因此 JAR 命名也随之简化在 Flink 1.15 之前则必须指定scala_version。未正式发布的自编译产物名称会带有SNAPSHOT后缀详见 加载文档中的编译说明。二、版本兼容性矩阵发布记录给出的官方兼容矩阵如下与 写入文档 和 读取文档 中的矩阵一致ConnectorFlinkStarRocksJavaScala1.2.151.16,1.17,1.18,1.19,1.202.1 and later82.11,2.121.2.141.16,1.17,1.18,1.19,1.202.1 and later82.11,2.121.2.121.16,1.17,1.18,1.19,1.202.1 and later82.11,2.121.2.111.15,1.16,1.17,1.18,1.19,1.202.1 and later82.11,2.12官方注意事项Flink connector 的最新版本一般只与最近三个 Flink 版本保持兼容。从矩阵可以读出两条实践结论使用 Flink 1.20 时1.2.11 及之后所有版本均可用而 Flink 1.15 只能停在 1.2.11 及更早版本所有列出的版本要求 StarRocks 2.1 及以上Java 8Scala 2.11/2.12。三、1.2.15多表事务与 exactly-once 残留事务加固发布时间2026 年 6 月 18 日。新特性支持多表事务 Stream LoadPR #487此前一次 Stream Load 事务只能针对单张目标表1.2.15 允许在一次事务中向多张表提交减少跨表写入时的事务数量。多表事务的用法与限制可参考仓库中的 多表事务加载文档。改进Merge Commit 支持记录数据质量错误日志PR #484启用 Merge Commit 后数据质量相关的错误信息会被记录便于排查脏数据。缺陷修复修复多表事务并发问题按表序列化加载、并对齐跨表提交PR #491对处于 PREPARE 状态的残留lingering事务当回滚失败时回退调用 FE 的 cancel API 进行清理PR #488——这是对 exactly-once 场景下作业异常退出后 StarRocks 侧残留事务清理链路的兜底构建列语句时不再给 DEFAULT 子句中的CURRENT_TIMESTAMP加引号PR #486避免生成非法建表语句。四、1.2.14Merge Commit 正式落地刷写窗口放宽到亚秒级发布时间2026 年 2 月 11 日。新特性支持 Merge CommitPR #474。这是该版本最重要的能力启用后多个 Flink sink 子任务的加载请求会在服务端合并进同一个 Stream Load 事务从而在不增加 StarRocks 事务数量的前提下提升sink.parallelism、获得更高吞吐。改进支持将sink.buffer-flush.interval-ms设置为小于 1 秒的值PR #475。这一变化在参数文档中有直接印证加载文档 中该参数说明写明早于 v1.2.14 的取值范围是 [1000, 3600000]v1.2.14 及以后为 (0, 3600000]即下限从 1 秒放宽到可配置亚秒级间隔利于低延迟场景支持配置事务 Publish 超时PR #480对应 sink 参数sink.publish-timeout.ms见 加载文档参数表。缺陷修复通过升级 guava 至32.0.1-jre修复 CVE-2023-2976PR #467。围绕 Merge Commit 的配套参数与调优1.2.14 的发布内容在 加载文档的 Merge Commit options 小节 中有完整定义值得在启用时一并掌握参数默认值说明sink.properties.enable_merge_commitfalse是否启用 Merge Commitsink.properties.merge_commit_interval_ms无启用时必填Merge 时间窗口窗口内的请求合并进同一事务sink.properties.merge_commit_parallel3每个 Merge Commit 事务在 StarRocks 侧创建的加载计划并行度区别于 Flink 侧的sink.parallelismsink.properties.merge_commit_asynctrue服务端返回模式异步模式配合 Flink checkpoint 提供 at-least-once 保证sink.merge-commit.max-concurrent-requestsInteger.MAX_VALUE单 sink 子任务并发 Stream Load 请求上限设为 0 可保证主键表顺序加载sink.merge-commit.chunk.size20971520单次 Stream Load 请求的 chunk 大小上限顺序模式0下默认改为 500 MB调优经验同样来自文档建议sink.buffer-flush.interval-ms不大于merge_commit_interval_ms使每个子任务在每个 Merge 窗口内至少刷写一次sink.buffer-flush.max-bytes应设为数倍于sink.merge-commit.chunk.size以允许至少累积满一个 chunk。五、1.2.12读端 warehouse 支持、安全策略与敏感日志脱敏发布时间2025 年 9 月 19 日。改进Source 支持指定 warehousePR #423读端可指向 Serverless 或计算组场景下的特定 warehouse新增安全策略PR #434错误日志敏感数据脱敏PR #446对应 sink 参数sink.sanitize-error-log默认false设为true后 connector 与 SDK 日志中的敏感行数据与列值会被打码参数说明同样记录在 加载文档 中标注Supported since 1.2.12支持为 Stream Load 事务接口配置prepared_timeoutPR #453对应参数sink.properties.prepared_timeout用于控制事务 prepare 阶段的等待时长文档同时指出Connector 1.2.12 且 StarRocks 3.5.4 时不设置该参数则回落到 FE 全局配置prepared_transaction_default_timeout_second默认 86400 秒。缺陷修复修复 source reader 在 open 失败时未被关闭的问题PR #441修复StreamLoadManagerV2.flush中任何异常导致误报成功的问题PR #451——该修复直接关系 exactly-once 语义的可靠性。六、1.2.11LZ4 压缩、Flink 1.20 适配与 JSON 数组包装选项发布时间2025 年 6 月 3 日。新特性CSV 格式支持 LZ4 压缩PR #408新增 Flink 1.20 支持PR #409。改进新增关闭将 JSON 包装成 JSON 数组的选项PR #344便于按 StarRocks 期望的单对象 JSON 格式刷写升级 FastJSON 修复 CVE-2022-25845PR #394从 warn 日志中移除数据行指标避免在日志中暴露数据负载PR #420。缺陷修复修复StarRocksDynamicTableSource影子克隆导致下推结果错误的问题修复后改用深拷贝PR #421。这是读端谓词下推正确性的一次关键修复。七、1.2.10读端复杂类型、异步事务接口与 socket 超时1.2.10 是读端能力的一次大版本。新特性支持读取 JSON 列PR #334支持读取 ARRAY、STRUCT、MAP 列PR #347JSON 格式 sink 支持 LZ4 压缩PR #354新增 Flink 1.19 支持PR #379。改进支持配置 socket 超时PR #319对应参数sink.socket.timeout-ms参数文档 标注Supported since 1.2.10默认-1表示不超时Stream Load 事务接口支持异步prepare与commitPR #328V2 事务接口sink.versionV2依赖 StarRocks 2.4 的事务接口在 1.2.10 起将 prepare/commit 异步化减少了 checkpoint 路径上的同步阻塞支持将 StarRocks 表的列子集映射到 Flink source 表PR #352使用 Stream Load 事务接口时支持指定 warehousePR #361。缺陷修复修复StarRocksDynamicLookupFunction中StarRocksSourceBeReader在读取完成后未关闭的问题PR #351修复加载空 JSON 字符串到 JSON 列时抛异常的问题PR #380。八、1.2.9对接 Flink CDC 3.0构建支持 schema 演进的流式 ELT1.2.9 的标志性变化是与 Flink CDC 3.0 集成便于从 MySQL、Kafka 等 CDC 源构建流式 ELT 管道并支持 schema change发布记录指向用户指南中的 Synchronize data with Flink CDC 3.0 章节对应 加载文档中的同名小节。新特性实现 catalog 以支持 Flink CDC 3.0PR #295实现 FLIP-191 新版 sink API 以支持 Flink CDC 3.0 的小文件合并诉求PR #301支持 Flink 1.18PR #305。缺陷修复修复误导性线程名与日志PR #290修复写多表时使用了错误的 stream-load-sdk 配置PR #298。九、1.2.8V1/V2 双版本 sink、exactly-once 实践基线与大量内部重构1.2.8 的发布说明强调两点支持 Flink 1.16/1.17当 sink 配置为 exactly-once 时建议设置sink.label-prefix以便 StarRocks 侧按 label 前缀识别并清理作业异常退出后残留的 PREPARE 状态事务加载文档的 Exactly Once 小节 对 1.2.8 同样有此推荐且明确 label 前缀需在全集群各类加载——Flink 作业、Routine Load、Broker Load——之间唯一。主要改进节选支持配置是否使用 Stream Load 事务接口保证 at-least-oncePR #228即sink.version的 AUTO/V1/V2 选择为 sink V1 增加重试指标PR #229重构将StarRocksSinkManagerV2、probeTransactionStreamLoad移入 stream-load-sdkPR #233、#240并支持 fastjson 替换为 jacksonPR #247根据 Flink 表 schema 自动识别 partial update不再强制用户显式写sink.properties.columnsPR #235支持处理update_before记录PR #250对应参数sink.ignore.update-before默认 true默认开启strip_outer_array与ignore_json_sizePR #259作业恢复且 sink 为 exactly-once 时尝试清理残留事务PR #271重试失败后返回首个异常PR #279。缺陷修复与测试修复StarRocksStreamLoadVisitor拼写问题PR #230与 fastjson classloader 泄漏PR #260新增从 Kafka 加载到 StarRocks 的测试框架PR #249并补充了 DataStream API 示例与 sink 文档PR #253、#262、#268、#275。十、版本选择建议与升级注意事项综合上述发布记录与仓库内文档可以给出如下选择依据按 Flink 版本反查 connectorFlink 1.20 用户可选 1.2.11 及以上Flink 1.15 用户最高只能使用 1.2.11。最新版一般只维护最近三个 Flink 版本升级 Flink 前先核对矩阵。exactly-once 场景建议使用 1.2.8 并配置sink.label-prefix1.2.10 的事务接口异步 prepare/commit 与 1.2.12 对 flush 误报成功的修复进一步提升了 exactly-once 链路的可靠性若使用 StarRocks 3.5.4可用sink.properties.prepared_timeout显式控制 prepare 超时。高并发写入1.2.14 可启用 Merge Commit将多个子任务请求合并为单事务配合上文的merge_commit_interval_ms、merge-commit.chunk.size、sink.merge-commit.max-concurrent-requests调优注意 Merge Commit 只保证 at-least-once不能与sink.semanticexactly-once同时使用。多表写入与跨表一致性使用 1.2.15其多表事务 Stream Load 及并发修复#487、#491可显著降低事务开销用法见 多表事务加载文档。读端Flink 消费 StarRocks需要读取 JSON/ARRAY/STRUCT/MAP 列或做列下推时至少使用 1.2.10读端参数与类型映射见 读取文档。参考文档Flink Connector 发布记录本文主体通过 Flink connector 持续写入 StarRockssink 参数全集、Exactly Once、Merge Commit 调优通过 Flink connector 读取 StarRocks读端参数与类型映射多表事务加载Stream Load 事务接口STREAM LOAD 语句与 Merge Commit 参数【免费下载链接】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),仅供参考