使用 Apache Flink® 通过多表事务加载数据
StarRocks Flink Connector 支持多表事务,可将 Flink 中的数据原子性地加载到多个表中。
使用场景
当单个 Flink 作业在一个处理周期内向同一 StarRocks 数据库中的多个表写入数据时,启用多表事务可保证:
- 跨表原子提交:在同一提交周期内写入不同表的数据以原子方式可见——要么全部成功,要么全部失败。
- 源事务完整性:完整的上游事务(例如来自 Kafka 的事务)不会被拆分到两个 StarRocks 事务中。
- 亚秒级数据新鲜度:数据通过
/api/transaction/load持续流入 StarRocks,并按sink.buffer-flush.interval-ms配置的间隔进行提交。
典型场景:
- 同步写入汇总表和明细表(例如
orders和order_items) - 将事件路由到不同的分区表(例如
events_202601、events_202602) - 单个作业维护多个相互关联的下游结果表
前提条件
要启用多表事务,您必须在 StarRocks v4.0 及以上版本(支持多表事务 Stream Load)上运行集群,并使用 v1.2.9 及以上版本的 StarRocks Flink Connector。
核心能力
| 能力 | 描述 |
|---|---|
| 跨表原子提交 | 同一刷新周期内的所有表共享一个 StarRocks 事务标签,Prepare 和 Commit 操作统一执行。 |
| 源事务完整性 | 提交时机由 transactionEnd 标志控制,仅在完整的源事务边界处进行提交。 |
| 亚秒级数据可见性 | 数据定期刷新到 StarRocks(/api/transaction/load),当满足 transactionEnd 和定时器条件时进行提交。 |
| N:1 事务映射 | 多个源事务可以在单个 StarRocks 事务中累积,无需按 1:1 映射。 |
| 分区内有序性 | keyBy(sourcePartition) 确保来自同一分区的事务在同一 sink 子任务中按顺序处理。 |
配置项
多表事务配置
sink.transaction.multi-table.enabled
- 类型:Boolean
- 默认值:
false - 描述:是否启用多表原子事务模式。
sink.transaction.multi-table.buffer-size
- 类型:Long
- 默认值:
134217728(128 MB) - 单位:字节
- 描述:多表事务模式下的全局缓冲区大小(字节)。当所有表的缓冲数据总量达到此阈值时,触发刷新。
加载相关配置
sink.version
- 推荐值:
V2 - 描述:必填项。
V1不支持事务 Stream Load 接口。
sink.semantic
- 推荐值:
at-least-once - 描述:多表模式当前仅支持
at-least-once。
database-name
- 推荐值:
* - 描述:通配符,用于启用动态多表路由。