「我的数据空间」实时计算实践笔记 · Flink SQL 系列本文档介绍使用 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 :-- 开启flink SQL hint options.SETtable.dynamic-table-options.enabledtrue;-- 指定 start-snapshot-id 读取UnBounded数据SELECT*FROMsample/* OPTIONS(streamingtrue, monitor-interval1s, start-snapshot-id3821550127947089987)*/;-- 指定 start-snapshot-id 和 end-snapshot-id 读取Bounded数据SELECT*FROMsample/* OPTIONS(streamingfalse, start-snapshot-id3821550127947089987, end-snapshot-id50127947089987382155)*/;在上述示例中我们使用streamingfalse默认为false即有界来表示该数据源是有界数据源通过start-snapshot-id和end-snapshot-id来界定消费区间。以下是各配置下的行为streamingstart-snapshot-idend-snapshot-id行为falsenullnull默认消费当前最新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 中使用如下配置来设置消费模式-- 指定消费模式为流模式 (默认消费模式)SETexecution.runtime-modestreaming;-- 指定消费模式为批模式SETexecution.runtime-modebatch;以下是消费模式和数据源配置下的行为execution.runtime-mode执行模式streaming 数据源有界无界行为streamingfalse默认行为流式消费区间段数据消费完成后作业置为FINISHED状态streamingtrue流式消费无界数据作业一直运行batchfalse批式消费一段数据消费完成后置为FINISHEDbatchtrue报错批模式无法消费无界数据源流式消费和增量消费的区别流式消费一般认为流式消费是流模式消费无界数据源正如上面章节表述流式消费无界数据源时首先使用最新 snapshot 作为 end-snapshot-id 消费全量数据然后周期性监控 snapshot 进行增量的读取。增量消费增量消费为消费一段数据即(start-snapshot-id, end-snapshot-id] 区间的数据前开后闭。当start-snapshot-id不配置时为全量消费当配置了 start-snapshot-id 则只消费区间内的新增数据。Flink 流批一体消费数据湖的建议当数据湖历史数据较多时建议使用 Flink batch模式消费历史数据创建 Flink SQL batch 作业insertintotableSELECT*FROMsample/* OPTIONS(streamingfalse, end-snapshot-id50127947089987382155)*/;跑批作业一定要使用 Flink SQL batch 作业当消费完成后再配置start-snapshot-id进行流式消费insertintotableSELECT*FROMsample/* OPTIONS(streamingtrue, start-snapshot-id50127947089987382155)*/;注意流式消费作业的start-snapshot-id必须与批作业的end-snapshot-id保持一致。该方法同样也可以用来回溯数据。iceberg v1 、v2表下Source算子并发度的配置在 iceberg v1 表和 v2 表在有界无界源下Source算子并发度有所区别iceberg表版本streaming模式行为v1false消费v1表有界源table.exec.iceberg.infer-source-parallelismtrue时默认值 并发度* Max(Min(**table.exec.iceberg.infer-source-parallelism.max100*, splitNumber) , 1)。splitNumber根据数据量的大小确定。v1truemonitor算子并发度为1。 reader算子如下计算 1.table.exec.resource.default-parallelism-1优先 2. 作业配置中Parallelism配置的并发度v2false指定start-snapshot-id monitor算子、 reader算子并发度皆为1。保证数据有序 不指定 start-snapshot-id 按照v1和false情况计算。 全量数据主键唯一不需要保证有序v2truemonitor算子、 reader算子并发度皆为1。保证数据有序iceberg v2表的增量读取iceberg v2表比较特殊由于存在更新需要消费changelog, 并且需要保证同一个key下的数据有序。iceberg 表在创建分区时也有一定的要求要求分区为主键的子集这是因为在读取时按照 partition 进行 merge-on-read 读取需要保证delete记录和insert记录在同一个分区下。所以在使用v2表作为source进行流式消费编程时要格外的小心尽量不要使用rebalance等操作 因为该操作会将原本有序的记录分发到不同的TM上进行指定在写入时无法保证有序。最好可以对 partition进行 keyby (groupby) 操作。本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。产品介绍见我的数据空间官网,支持私有化部署与 OEM 合作。