主题
FlinkSQL + Iceberg 搭实时看板:Append 流与 Changelog 流是两套写法
Flink作为一款优秀的实时计算引擎,在稳定性、容错、实时性多个方面都有着相当不错的表现,对SQL的兼容与支持让Flink的能力大幅提升,简化了接入流程、降低了学习难度。本文主要介绍如何用 FlinkSQL 把实时聚合结果写进数据湖,再由 BI 工具搭建实时看板。 看板侧一般支持 MySQL、Doris、Hive、Iceberg 这几种主流数据源。综合考虑数据量、查询性能与功能丰富程度,针对数据量不是太大、对查询性能有较大需求的同学,可以选择将数据写入Doris,数据量较大,对查询速度要求不高的同学,可以将数据写入Iceberg。下文将以Iceberg为底座来演示,以Doris为底层数据存储同理。 实现实时看板目前存在着两种方式:
- 计算逻辑放在FlinkSQL中,使用FlinkSQL实时计算进行加速,将处理后的数据直接写入下游Iceberg或Doris表中。
- 计算逻辑下沉到底层数据库中,FlinkSQL只进行简单的逻辑处理,保存明细数据到下游表中,在看板中配置不同的查询聚合逻辑。 方式1和方式2都有各自的特点,方式1的特点在于计算逻辑在FlinkSQL中,下游表的数据都是直接聚合好的结果,对于查询来说较为友好,不需要再做额外的处理,但相对而言指标维度已经固定,无法再进行维度下钻。方式2的特点在于数据维度较全,可以根据需求随时更改,方便进行下钻分析。对于方式1和方式2的选择需要基于业务逻辑,在实现时可以根据具体需要将部分不需要的维度在Flink作业中提前聚合,在Flink中提前将数据打宽,减少看板每次查询时的消耗。 我们以销售订单为例,其中订单存在以下维度信息:
sql
CREATE TABLE kafka_catalog.default.order_tb(
id BIGINT COMMENT 订单ID, -- primary key
goods_type VARCHAR COMMENT 物品类型,
goods_cnt BIGINT COMMENT 物品数量,
shop_name VARCHAR COMMENT 门店名, -- primary key
region_name VARCHAR COMMENT 门店区域,
pay_time TIMESTAMP COMMENT 订单付款时间,
up_time TIMESTAMP COMMENT 订单更新时间
)针对上游数据的不同,我们将内容分为两大块,分别是针对上游是Append数据流,以及上游是Changelog数据流。 Append流可以认为都是INSERT的数据,通常来自用户自身业务的打点数据亦或来自于监控打点的数据,不存在变更和主键关系。对于MySQL等带主键的表进行Batch读,得到的即是一条Append流,对Iceberg v1、Hive表进行Batch/增量读,得到的也是Append流。 Changelog流是存在增删改的数据,通常由业务MySQL、TiDB库变更日志同步得到。通过实时集成从MySQL或者TiDB等数据表得到的数据流即是Changelog流,对Iceberg V2表的增量读得到的也是Changelog流,通常Changelog流中的字段会有主键存在。
建表
首先,需要先在平台空间下建一张Iceberg表来存储聚合后的数据,因为涉及到计算逻辑,上游Append流经过聚合计算后会产生changelog数据,为了更新能够同步到Iceberg表中,在这里我们选择V2表。
假设我们仅需要关注门店商品的销售情况,时间粒度仅需要到小时,那么新建的Iceberg表需要用到如下字段,去掉订单Id和订单更新时间(append流订单并不会更新),将主键设为goods_type和shop_name,同时设置表按shop_name进行分区。
sql
CREATE TABLE iceberg_catalog.iceberg.order_tb(
goods_type VARCHAR COMMENT 物品类型, -- primary key
goods_cnt BIGINT COMMENT 物品数量,
shop_name VARCHAR COMMENT 门店名, -- primary key
region_name VARCHAR COMMENT 门店区域,
pay_time TIMESTAMP COMMENT 订单付款时间
)
数据同步
Append 流
Append流指的是上游输入的数据都是INSERT类型的,不包含UPDATE和DELETE类型的数据。假设我们 order_tb表中的数据由线上服务手动打点到Kafka中,不存在更新和删除情况,那么Flink消费到的就是一条Append流。在Append流中,可以根据需求做相应处理。 在进行了建表后,就可以根据按照门店和小时聚合的逻辑来写SQL。按小时聚合的话,会导致上游数据一小时才会下发一次,不利于看板的实时化,针对这样的问题,在1.13版本之前,可以先使用窗口的early-fire,在1.13+版本可以使用CUMULATE WINDOW. 由于Iceberg支持upsert写入,那么后续early-fire的值能够覆盖前面的值,能够保证数据的准确。
sql
-- 设置checkpoint间隔,如果数据量不大可以调小至1min
SET execution.checkpointing.interval=3minute;
SET execution.checkpointing.timeout=30minute;
SET execution.checkpointing.mode=EXACTLY_ONCE;
-- 写V2表一定要设置,否则可能会导致下游写重复
SET execution.checkpointing.tolerable-failed-checkpoints=0;
-- 开启early-fire并设置时间为3分钟下发一次
SET table.exec.emit.early-fire.enabled=true;
SET table.exec.emit.early-fire.delay=3min;
CREATE TABLE source_tb (
pay_time_timestamp as TO_TIMESTAMP(pay_time_str),
pay_time_timestamp as TO_TIMESTAMP_LTZ(pay_time_long),
pay_time_timestamp as TO_TIMESTAMP_LTZ(cast(pay_time_long_str as bigint)),
pay_time_timestamp as cast(pay_time_timestamp as timestamp(3))
WATERMARK FOR pay_time AS pay_time_timestamp - INTERVAL '5' SECOND
) LIKE kafka_catalog.default.order_tb;
INSERT INTO iceberg_catalog.iceberg.order_tb
SELECT goods_type, sum(goods_cnt) AS cnt, shop_name, region_name, TUMBLE_END(pay_time, INTERVAL '1' HOUR) AS pay_time
FROM source_tb
GROUP BY TUMBLE(pay_time, INTERVAL '1' HOUR), goods_type, shop_name, region_name;Changelog 流
如果上游数据是从MySQL binlog中收集得出,那么就存在UPDATE和DELETE类型的数据,针对这种changelog流,目前FlinkSQL中的窗口并不能支持,需要借助groupAgg来实现窗口的逻辑。同样,建立和前面一样的Iceberg表,并使用如下的逻辑进行处理。
sql
-- 设置checkpoint间隔,如果数据量不大可以调小至1min
SET execution.checkpointing.interval=3minute;
SET execution.checkpointing.timeout=30minute;
SET execution.checkpointing.mode=EXACTLY_ONCE;
-- 写V2表一定要设置,否则可能会导致下游写重复
SET execution.checkpointing.tolerable-failed-checkpoints=0;
INSERT INTO iceberg_catalog.iceberg.order_tb
WITH order_view AS (
SELECT goods_type, goods_cnt, shop_name, region_name,
DATE_FORMAT(pay_time, 'yyyy-dd-MM hh') AS new_time
FROM kafka_catalog.default.order_tb
)
SELECT goods_type, sum(goods_cnt) AS cnt, shop_name,
region_name, TO_TIMESTAMP(new_time, 'yyyy-dd-MM hh')
FROM order_view
GROUP BY goods_type, shop_name, region_name, new_time;将时间列pay_time粒度退化到小时,再用groupAgg进行聚合,求门店某种商品的在一小时内销售的总数。通过groupAgg的处理后,数据会被实时更新,有新数据来会先发一条update_before数据,再下发一条update_after数据来进行变更,因为Iceberg支持upsert写入,如果数据频繁回撤导致下游抖动,可以只把 update_after 写下去、丢掉 update_before,写入量减半。但要清楚这是拿正确性换平滑:一旦下游不是按主键覆盖,丢掉的前镜像就补不回来了。
看板搭建
数据写入后,就可以在 BI 工具里搭建实时看板了。步骤上通常是三步:接入数据源(选到 Iceberg 表所在的 catalog/库)→ 建一个数据模型(在模型里用 SQL 把需要的字段取出来,并把度量列标成指标)→ 基于模型建图表。
需要确认一件事:所选 BI 工具是否支持 Iceberg 作为查询底座。部分工具只支持 Hive/MySQL/Doris,这种情况下把 Flink 的落点换成 Doris 即可,上游 SQL 不用改。
各门店销量看板
基于上面的数据模型创建一个直方图看板:维度选 pay_time,指标选每一类商品的销售数量。选中某个小时的数据即可继续做维度下钻,查看该小时内具体门店的销量。
门店分时总销量看板
换成以门店为维度、分时段汇总销量的看板,同一份数据模型即可复用。
由于使用Flink实时写入下游Iceberg表,下游表中的数据会实时更新,看板可以根据需要设置使用时更新或者定时更新。
性能优化:
- 提升查询速度
- 由于看板的每一次查询都是全量的,SQL逻辑越复杂,涉及的表越多,就会导致查询的速度越慢。针对部分数据量大、涉及多表的场景,建议提前在Flink中处理好数据,使用Flink的Join能力,将数据打宽和计算逻辑提前到Flink中。
- 提升Flink消费Iceberg时效性 由于Iceberg并不像消息队列一样,能够稳定在秒级延迟,目前仅可以稳定在分钟级延迟,但显然1分钟和30分钟,也都算分钟延迟,通过下面几点可以尽可能提升数据流转的速度和时效性。
- 控制好上游数据可见延迟 首先我们关注上游数据的产出速度,如果是类似于Kafka这一类的实时系统,那么这一环可以忽略,如果上游是Hive或Iceberg,那么需要注意一下数据可见性问题,对于Iceberg系统来说,数据的commit依赖上游作业的checkpoint,想要提升看板的实时性,首先需要降低source表上游作业的checkpoint间隔,2~10min一次为宜。
- 控制好同步作业的checkpoint间隔 与上面同理,实时看板依赖的iceberg表也依赖Flink作业commit,我们也需要保证尽可能短的checkpoint间隔,同样也以2~10min为宜。
- 开启mini batch,设置好Source表的参数 Flink作业可以通过开启mini batch来增加执行效率
sql
set table.exec.mini-batch.enabled=true;
set table.exec.mini-batch.size=100;
set table.exec.mini-batch.allow-latency=5s;如果Source表也是Iceberg表,那么需要设置好monitor的间隔和每次下发的split数据量
sql
set table.exec.source.cdc-events-duplicate=true; -- 因为iceberg目前支持upsert读,加上去重修复求sum异常
select * from iceberg_catalog.iceberg.iceberg_source
/*+OPTIONS('streaming'='true', 'monitor-interval'='30s', 'max-planning-snapshot-count'='5') */目前增量消费 Iceberg 的速度由 monitor 间隔和单轮规划的 snapshot 数控制,不同的作业有着不一定的处理能力,需要手动调整这两个参数。
本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览与技术博客,交流合作 QQ:1559851993。