主题
Flink 流批一体读 Iceberg:snapshot 区间的参数矩阵与 V2 表的读取约束
本文档介绍使用 Flink 流批模式下消费数据湖的一些关键点和细节。
使用 iceberg + Flink 构建实时链路
基于 iceberg 的快照(snapshot)特性及 Flink 流式消费的能力,我们可以将链路延迟降低到分钟级别。 链路的时效性取决于上游 Flink 作业的 checkpoint 间隔,每次 checkpoint 完成后都会形成一个新的快照,该快照包含了在该 checkpoint 间隔内写入的所有数据。
每个Snapshot 都是保存了当时时刻的全局数据,即表的所有数据。这是通过引用所有的 manifest 文件(manifest文件 是所有的数据文件列表集合,用来记录数据文件和 partition 的对应关系)来实现的。本次 snapshot 内新增的 manifest 文件会被标记为 added, 当增量读取时,只需要读取 added 的manifest 即可读取到本次 snapshot 新增的数据。
下游 Flink 作业流式消费时会周期性的扫描最新的快照(latest-snapshot), 并读取 (last-snapshot, least-snapshot] 区间内所有的新增数据(增量读)。

有界数据源、无界数据源
Flink 流式消费与之前的批模式不同的是它的源是无界数据源(UnBounded) , 上游数据会一直到来,因此该作业会一直运行。对应到数据湖中,当上游 iceberg 表不断有新的 snapshot 生成,可以理解为一个无界数据源。当我们指定起始 snapshot (start-snapshot-id) 和终止 snapshot (end-snapshot-id) 时可以认为是一个有界源。 使用 Flink 读取 iceberg :
sql
-- 开启flink SQL hint options.
SET table.dynamic-table-options.enabled=true;
-- 指定 start-snapshot-id 读取UnBounded数据
SELECT * FROM sample /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s', 'start-snapshot-id'='3821550127947089987')*/ ;
-- 指定 start-snapshot-id 和 end-snapshot-id 读取Bounded数据
SELECT * FROM sample /*+ OPTIONS('streaming'='false', 'start-snapshot-id'='3821550127947089987', 'end-snapshot-id'='50127947089987382155')*/ ;在上述示例中我们使用 streaming=false (默认为false,即有界)来表示该数据源是有界数据源,通过start-snapshot-id 和 end-snapshot-id 来界定消费区间。 以下是各配置下的行为:
| streaming | start-snapshot-id | end-snapshot-id | 行为 |
|---|---|---|---|
| false | null | null | 默认消费当前最新snapshot的全量数据 |
| false | 不指定 | 指定 | 消费(-1, end-snapshot-id] 区间的全量数据 |
| false | 指定 | 指定 | 消费(start-snapshot-id, end-snapshot-id]区间的新增数据 |
| false | 指定 | 不指定 | 报错,不是有界源 |
| true | 不指定 | 不指定 | 无界源, 只能使用流式消费模式,首先消费(-1, least-snapshot]消费全量数据,然后(last-snapshot, least-snapshot] 增量消费 |
| true | 指定 | 不指定 | 无界源,只消费(start-snapshot-id, latest-snapshot] 增量数据 |
| true | - | 指定 | 报错, 不是无界源 |
Flink 流批消费模式
数据源的有界无界和 Flink 的消费模式没有关联,即可以使用流模式来消费无界数据源,也可以使用流模式来消费有界数据源。 在 Flink 中使用如下配置来设置消费模式:
sql
-- 指定消费模式为流模式 (默认消费模式)
SET execution.runtime-mode = streaming;
-- 指定消费模式为批模式
SET execution.runtime-mode = batch;以下是消费模式和数据源配置下的行为:
| execution.runtime-mode(执行模式) | streaming (数据源有界无界) | 行为 |
|---|---|---|
| streaming | false | 默认行为,流式消费区间段数据,消费完成后作业置为FINISHED状态 |
| streaming | true | 流式消费无界数据,作业一直运行 |
| batch | false | 批式消费一段数据,消费完成后置为FINISHED |
| batch | true | 报错,批模式无法消费无界数据源 |
流式消费和增量消费的区别
流式消费:一般认为流式消费是流模式消费无界数据源,正如上面章节表述,流式消费无界数据源时,首先使用最新 snapshot 作为 end-snapshot-id 消费全量数据,然后周期性监控 snapshot , 进行增量的读取。 增量消费:增量消费为消费一段数据,即(start-snapshot-id, end-snapshot-id] 区间的数据,前开后闭。当start-snapshot-id不配置时为全量消费,当配置了 start-snapshot-id , 则只消费区间内的新增数据。
Flink 流批一体消费数据湖的建议
当数据湖历史数据较多时,建议使用 Flink batch模式消费历史数据,创建 Flink SQL batch 作业
sql
insert into table
SELECT * FROM sample /*+ OPTIONS('streaming'='false', 'end-snapshot-id'='50127947089987382155')*/ ;跑批作业一定要使用 Flink SQL batch 作业
当消费完成后,再配置start-snapshot-id进行流式消费:
sql
insert into table
SELECT * FROM sample /*+ OPTIONS('streaming'='true', 'start-snapshot-id'='50127947089987382155')*/ ;注意流式消费作业的start-snapshot-id必须与批作业的end-snapshot-id保持一致。 该方法同样也可以用来回溯数据。
iceberg v1 、v2表下Source算子并发度的配置
在 iceberg v1 表和 v2 表在有界无界源下Source算子并发度有所区别:
| iceberg表版本 | streaming模式 | 行为 |
|---|---|---|
| v1 | false | 消费v1表有界源:table.exec.iceberg.infer-source-parallelism=true 时(默认值), 并发度=* Max(Min(**table.exec.iceberg.infer-source-parallelism.max=100* , splitNumber) , 1) 。splitNumber根据数据量的大小确定。 |
| v1 | true | monitor算子并发度为1。 reader算子如下计算: 1. table.exec.resource.default-parallelism=-1 (优先) 2. 作业配置中Parallelism配置的并发度 |
| v2 | false | 指定start-snapshot-id : monitor算子、 reader算子并发度皆为1。(保证数据有序) 不指定 start-snapshot-id: 按照v1和false情况计算。 (全量数据,主键唯一,不需要保证有序) |
| v2 | true | monitor算子、 reader算子并发度皆为1。(保证数据有序) |
iceberg v2表的增量读取
iceberg v2表比较特殊,由于存在更新,需要消费changelog, 并且需要保证同一个key下的数据有序。iceberg 表在创建分区时也有一定的要求,要求分区为主键的子集,这是因为在读取时按照 partition 进行 merge-on-read 读取,需要保证delete记录和insert记录在同一个分区下。
所以在使用v2表作为source进行流式消费编程时要格外的小心,尽量不要使用rebalance等操作, 因为该操作会将原本有序的记录分发到不同的TM上进行指定,在写入时无法保证有序。最好可以对 partition进行 keyby (groupby) 操作。
本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览与技术博客,交流合作 QQ:1559851993。