使用 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.transaction.multi-table.mini-switch-interval-ms
- Type: Long
- Default:
-1(自动) - Unit: ms
- Description: 每个分区执行 Chunk 切换之间的最小间隔(毫秒)。在此间隔内,源事务会被批量合并到一次 Stream Load 中,从而限制 HTTP 请求数量。
-1表示自动计算为min(1000, max(500, sink.buffer-flush.interval-ms/4));设置正数时则使用指定值。
sink.transaction.multi-table.min-switch-bytes
- Type: Long
- Default:
1048576(1 MB) - Unit: Bytes
- Description: 基于时间间隔触发的 Chunk 切换仅在某个 Region 的活动 Chunk 至少累积了此字节数后才会执行。这样,低流量分区可以将更多源事务批量合并到一次 Stream Load 中,而不是产生大量微小请求。达到半满时的预留空间、缓冲区大小导致的内存压力,或者自该分区上次切换后经过完整的
sink.buffer-flush.interval-ms,都会强制执行切换,不受此大小阈值限制。因此,即使设置了较大的阈值,持续的低流量数据流仍会大约按照 flush 间隔进行切换,不会被无限期地延迟。设置为<=0可禁用大小限制,恢复之前按时间间隔切换的行为。
sink.transaction.multi-table.max-txn-bytes
- Type: Long
- Default:
0(无限制) - Unit: Bytes
- Description: 在等待
txnEnd期间,单个 sink 子任务可以缓冲的正在进行中的源事务数据的硬上限(字节)。对于在txnEnd之前无法 flush 的数据,Writer 不会阻塞(参见限制 7),因此单个源事务可能增长到超过2 × buffer-size。此选项用于限制这种增长;超过 限制时,任务将因明确的错误而失败。0表示不限制大小,此时唯一的限制是 TaskManager 的堆内存。
加载相关配置
sink.version
- 推荐值:
V2 - 描述:必填项。
V1不支持事务 Stream Load 接口。
sink.semantic
- 推荐值:
at-least-once - 描述:多表模式当前仅支持
at-least-once。
database-name
- 推荐值:
* - 描述:通配符,用于启用动态多表路由。
table-name
- 推荐值:
* - 描述:通配符,用于启用动态多表路由。
sink.buffer-flush.interval-ms
- 推荐值:
1000 - 描述:控制提交周期。可将其设置为
1000以实现约一秒的数据新鲜度。
sink.properties.format
- 推荐值:
json - 描述:数据格式。
sink.properties.strip_outer_array
- 推荐值:
true - 描述:是否去除最外层的数组结构。
接口
StarRocksRowData
public interface StarRocksRowData {
String getUniqueKey(); // Region routing key (nullable; auto-derived from database.table)
String getDatabase(); // Target database
String getTable(); // Target table
String getRow(); // Row data in JSON format
/**
* Indicates this is the last row of a source transaction batch.
* Used by multi-table transaction mode to determine safe commit points:
* the connector only commits when the most recent write had this flag set,
* ensuring no partial source transaction is committed.
*/
default boolean isTransactionEnd() {
return false;
}
/**
* Returns the source partition ID for this row.
* Used by multi-table transaction mode to track per-partition transaction
* boundaries. Returns -1 when partition tracking is not applicable.
*/
default int getSourcePartition() {
return -1;
}
}
DefaultStarRocksRowData
public class DefaultStarRocksRowData implements StarRocksRowData {
// 基本字段
private String uniqueKey;
private String database;
private String table;
private String row;
// 多表事务字段
private boolean transactionEnd; // Source transaction end marker
private int sourcePartition = -1; // Source partition ID (for keyBy ordering)
// 构造函数
public DefaultStarRocksRowData();
public DefaultStarRocksRowData(String database, String table);
public DefaultStarRocksRowData(String uniqueKey, String database, String table, String row);
// Setter 方法
public void setUniqueKey(String uniqueKey);
public void setDatabase(String database);
public void setTable(String table);
public void setRow(String row);
public void setTransactionEnd(boolean transactionEnd);
public void setSourcePartition(int sourcePartition);
// Getter 方法(继承自 StarRocksRowData)
public String getUniqueKey();
public String getDatabase();
public String getTable();
public String getRow();
public boolean isTransactionEnd();
public int getSourcePartition();
}
用户实现的组件
用户需要实现一个 KeyedProcessFunction(在本文档中称为 TransactionAssembler),它:
- 按源分区作为键并在事务中缓冲数据行
- 仅在源事务关闭时(例如,收到
TXN_END时)才发出所有行 - 在最后一行设置
transactionEnd=true - 在每一行上设置
sourcePartition
无需自定义 SinkFunction — 标准连接器 API(SinkFunctionFactory.createSinkFunction())可处理一切。
完整示例
StarRocks 表 DDL
CREATE DATABASE `test`;
CREATE TABLE `test`.`orders` (
`order_id` BIGINT NOT NULL,
`customer_id` BIGINT NOT NULL,
`total_amount` DECIMAL(10,2) DEFAULT "0",
`order_status` VARCHAR(32) DEFAULT ""
) ENGINE=OLAP PRIMARY KEY(`order_id`)
DISTRIBUTED BY HASH(`order_id`)
PROPERTIES("replication_num" = "1");
CREATE TABLE `test`.`order_items` (
`item_id` BIGINT NOT NULL,
`order_id` BIGINT NOT NULL,
`product_name` VARCHAR(128) DEFAULT "",
`quantity` INT DEFAULT "0",
`price` DECIMAL(10,2) DEFAULT "0"
) ENGINE=OLAP PRIMARY KEY(`item_id`)
DISTRIBUTED BY HASH(`item_id`)
PROPERTIES("replication_num" = "1");