这个话题我得先泼一盆冷水:很多人以为Flink SQL就是把离线SQL改个引擎就能跑实时,结果一上手就发现,source定义不对、状态没清理、窗口没触发,问题一个接一个。我在好几个实时数仓项目里都见过这种场面。这篇文章不扯虚的,把我用Flink SQL做实时同步和实时数仓的真实经验、踩过的坑、调优思路全部整理出来,看完你至少能少走三个月的弯路。
先说清楚一件事:Flink SQL不是来替代DataStream API的,它是为了解决“实时计算开发门槛太高”这个痛点而存在的。一个团队里,Java工程师写DataStream逻辑可能要几天,但同样是这批人,用Flink SQL写一套从Kafka到ClickHouse的实时ETL,一个下午就能跑通。这不是说DataStream没用,而是SQL方式让更多没有深度流计算背景的人也能参与实时数仓建设,这才是它最大的价值。对于想进入大数据领域的人、正在做实时数仓的工程师、以及那些被数据同步任务折腾得头疼的运维同学,这篇文章都值得你读完。
1. 为什么实时场景里我首选Flink SQL
1.1 离线和实时终于能用同一种语言了
过去做实时计算,逻辑复杂点就得写MapFunction、RichFlatMapFunction,还要自己管状态、管watermark、管exactly-once,一套流程下来代码量巨大。Flink SQL最大的意义在于,你写一条CREATE TABLE和INSERT INTO,底层那些复杂机制框架全包了。
我用一个生活化的例子来解释:DataStream API像是手动挡开车,每个操作都要自己换挡、踩离合,动力性能上限高;Flink SQL像是自动挡,你只管油门刹车,换挡逻辑交给变速箱。日常通勤没人会天天想手动换挡的事,Flink SQL就是绝大多数实时计算场景里那辆自动挡的车。
从开发效率上说,一个实时指标需求,用DataStream API从编码到调试可能要两天,Flink SQL通常半天内能交活。同一个团队同时维护两套技术栈,成本很高,所以现在的项目里我们默认用SQL,遇到SQL表达不了的特殊场景(比如某些自定义UDF都不方便实现的算子),才退回去用DataStream兜底。
1.2 流批一体不是营销概念,是真实收益
很多公司有两条计算链路:离线用Hive/Spark跑T+1报表,实时用Flink跑分钟级指标。结果同一个指标,离线算出来一个数,实时算出来另一个数,业务方天天找你对口径。Flink SQL的流批一体特性解决的就是这个问题:同一套SQL逻辑,既能跑流式任务处理实时数据,也能跑批式任务处理历史数据,因为引擎对两套执行做了统一优化。
在1.17之后,Flink SQL的流批模式切换只需要配置execution.runtime-mode,一套SQL两种模式跑,口径天然对齐。我们曾经把一个订单汇总指标从双链路改成Flink SQL单链路,离线批处理每天凌晨跑一次,实时流处理每5分钟出一个结果,两边数据完全对得上。这件事做成了,整个数据团队的口径扯皮问题都少了一半。
1.3 生态位优势太明显
Flink SQL能活这么好,很大程度靠的是连接器生态。Kafka、MySQL、PostgreSQL、ClickHouse、Hudi、Iceberg、Doris,主流的数据源和数据湖都有官方或社区连接器。你在SQL里写个WITH连接器参数,读写就通了,比原来写个自定义Sink省事太多。
更关键的是,Flink SQL的表结构能跟Hive Metastore打通。我们线上实时数仓直接复用离线数仓的表元数据,用Hive Catalog把Flink表关联到Hive的元数据服务上,实时链路和离线链路共用同一套表结构定义。可以说,Flink SQL已经成了不少公司实时数据架构的事实标准。
2. Flink SQL核心原理,不懂这些写不出稳的任务
2.1 动态表和连续查询是两根支柱
Flink SQL对流的抽象是动态表。流数据每来一条,动态表就多一行(或者更新一行)。你写的SELECT语句,其实是一条连续查询,它永远不会结束,会持续根据新到的数据更新自己的结果。这和离线SQL查询一个“静止的表”是完全不同的心智模型。
我见过很多新人把Flink SQL当作离线SQL来写,写完之后发现结果老是不对,就是因为没理解“查询是连续的”。比如一条简单的分组统计:
SELECT user_id, COUNT(*) AS cnt FROM orders GROUP BY user_id;离线里这跑完给个最终结果,任务就结束了。Flink SQL里这是无限流,每个user_id的cnt会不断更新,并且持续往下游输出。对下游来说,相当于一直收到“这个用户最新的累计值”。如果你下游是个打印或者Redis,就得做好覆盖写、upsert的准备,而不能像离线那样只跑一次。
2.2 时间属性决定窗口任务靠不靠谱
做实时计算必须处理时间,Flink SQL里最核心的就是声明时间属性。两种时间,处理时间就是机器当前时间,简单但结果不确定;事件时间是数据自带的时间戳,能处理乱序和延迟,但需要配置watermark来告诉引擎“我可以等多久”。
我强烈建议,只要业务允许,统统用事件时间。曾经有个交易大屏项目,最初图省事全用处理时间,结果上游一个批次数据因为网络抖动延迟了30秒到达,大屏指标瞬间跳变,运维被业务方点名批评。后来全部改造为事件时间,加上watermark和allowedLateness,数据恢复之后指标自动修正,再没人投诉过。
定义事件时间的标准写法:
CREATE TABLE orders ( order_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND ) WITH (...);WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND的意思是:引擎认为,比当前已见的最大事件时间晚5秒以上的数据都已经到齐了。这5秒是给乱序数据的缓冲。这个值不能拍脑袋定,要看上游Kafka的消息延迟分布。我们一般先在离线环境统计P95延迟,再定这个参数,太大则实时性差,太小则丢数据。
2.3 窗口函数是实时统计的弹药库
Flink SQL窗口分三种:滚动窗口、滑动窗口、会话窗口。滚动窗口最常用,比如每5分钟统计一次订单金额:
SELECT TUMBLE_START(order_time, INTERVAL '5' MINUTE) AS win_start, TUMBLE_END(order_time, INTERVAL '5' MINUTE) AS win_end, SUM(amount) AS total_amount FROM orders GROUP BY TUMBLE(order_time, INTERVAL '5' MINUTE);注意分组键是整一个TUMBLE窗口函数,不能只写时间字段。滑动窗口比如每5分钟统计一次、但窗口长度为15分钟,用来做平滑趋势:
SELECT HOP_START(order_time, INTERVAL '5' MINUTE', INTERVAL '15' MINUTE'), SUM(amount) FROM orders GROUP BY HOP(order_time, INTERVAL '5' MINUTE', INTERVAL '15' MINUTE');会话窗口则适合用户行为分析,比如用户连续两次操作间隔超过10分钟,就断开为一次会话。这里有个细节,千万别在SQL里用GROUP BY去重来模拟会话切割,状态会失控的,老老实用SESSION窗口函数。
窗口还有一种写法是OVER窗口,做流式累计计算(比如用户至今累计消费额)。它不按时间分组,而是按行或时间范围滑动的。OVER窗口在实时数仓里做累计快照非常有用,但注意它要基于主键去重后才能用,否则结果可能多条重复。
3. 实战记录:用Flink SQL把MySQL实时同步到ClickHouse
3.1 场景背景与整体链路设计
我们有一个订单系统在MySQL里,业务方需要一个实时的大屏看板,还要支持多维度的OLAP查询。MySQL显然扛不住大屏的高频聚合查询,所以目标是把订单表实时同步到ClickHouse。这个场景非常典型,基本是每个实时数仓项目都要做的事。
整体链路是:MySQL Binlog → Kafka → Flink SQL → ClickHouse。有人会问为啥不直接MySQL Binlog到Flink再到ClickHouse,非要中间插个Kafka?原因有两个:
一是削峰填谷,业务高峰时MySQL的Binlog量巨大,Flink任务重启或者做Savepoint恢复期间,Kafka能先把数据存住,不会丢。二是解耦,下游可以同时接多个Flink任务,比如一份数据既进ClickHouse做OLAP,又进Elasticsearch做搜索,Kafka作为数据总线非常方便。
3.2 三张核心建表语句与参数说明
第一步,在Flink SQL里建Kafka的source表。注意格式和元数据字段的保留:
CREATE TABLE order_cdc ( id BIGINT, order_no STRING, user_id BIGINT, amount DECIMAL(10,2), create_time TIMESTAMP(3), op_ts TIMESTAMP(3) METADATA FROM 'timestamp' -- 记录Binlog时间 ) WITH ( 'connector' = 'kafka', 'topic' = 'dwd_order_cdc', 'properties.bootstrap.servers' = 'kafka01:9092,kafka02:9092', 'properties.group.id' = 'flink-cdc-order-group', 'format' = 'debezium-json', 'scan.startup.mode' = 'earliest-offset' );debezium-json格式很讲究。Flink CDC(也就是Flink CDC连接器)直接捕获MySQL Binlog后,默认输出的是Debezium格式的消息,里面有before、after、op这些字段。用format = 'debezium-json',Flink能自动识别INSERT、UPDATE、DELETE事件,并转成对目标表的upsert/delete操作。如果用错了格式,比如用了默认json,你会看到所有Binlog事件都被当成INSERT,更新操作直接把旧数据再插一遍,ClickHouse里垃圾数据一大堆。
第二步,建ClickHouse结果表。ClickHouse连接器有两个关键参数:
CREATE TABLE order_ch ( id BIGINT, order_no STRING, user_id BIGINT, amount DECIMAL(10,2), create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'clickhouse', 'url' = 'clickhouse://clickhouse01:8123', 'database-name' = 'bi', 'table-name' = 'dwd_order_ch', 'table.catalog.name' = 'default', 'sink.batch.size' = '500', 'sink.batch.interval' = '2s', 'sink.max-retries' = '3' );PRIMARY KEY (id) NOT ENFORCED这个写法很多人看不懂。NOT ENFORCED的意思是:这个主键只是告诉Flink SQL“数据以id作为更新逻辑的键”,但Flink不负责检查主键唯一性,真正的主键约束由ClickHouse的表引擎保证。如果我们用PRIMARY KEY (id)不带NOT ENFORCED,那需要表上定义唯一约束,有些连接器会直接报错。
sink.batch.size和sink.batch.interval是ClickHouse异步批量写入的触发条件,满足任何一个就会把攒着的一批数据写出去。调小batch.size实时性更高,但ClickHouse写入频率太高会产生大量小parts,后台Merge压力变大;调大了ClickHouse性能好,但实时性下降。我们项目里压测下来,每秒几千条数据的量级,batch.size 500、interval 2秒是一个平衡点。
第三步,也是最容易被忽略的一步:主键去重。ClickHouse连接器虽然支持upsert,但Flink SQL需要明确知道哪一列是更新的依据。如果原始订单表里有重复的op_ts数据或者下游重复投递,直接写入会重复。我们通常先在SQL里做一次按主键的聚合去重:
CREATE VIEW order_dedup AS SELECT id, order_no, user_id, amount, create_time FROM ( SELECT id, order_no, user_id, amount, create_time, ROW_NUMBER() OVER (PARTITION BY id ORDER BY create_time DESC) AS rn FROM order_cdc ) WHERE rn = 1;这就是热搜里“sql语句去重”最常见的实时版解法。之后INSERT INTO order_ch时直接select这张视图。
3.3 启动任务与Checkpoint配置
建完表,提交任务的SQL就一句话:
INSERT INTO order_ch SELECT id, order_no, user_id, amount, create_time FROM order_dedup;但别高兴太早,提交之前必须检查Checkpoint配置,否则任务跑几天就会出现状态越来越大的问题。提交yarn-session任务时,我习惯在Flink SQL的SET命令里把这些参数写清楚:
SET execution.checkpointing.interval = 60s; SET execution.checkpointing.mode = EXACTLY_ONCE; SET execution.checkpointing.timeout = 30s; SET state.backend.type = rocksdb; SET state.checkpoint-storage = filesystem; SET state.checkpoints.dir = hdfs:///flink/checkpoints;Checkpoint是Flink SQL的命根子。如果不开Checkpoint,任务重启后状态全丢,去重逻辑失效,下游会涌入大量重复数据。用RocksDB状态后端是因为大状态场景下,Java堆根本装不下,RocksDB把状态放磁盘,用内存做缓存,虽然单次访问慢一些,但容量大得多。
关于增量Checkpoint,如果是RocksDB,默认就已经启用了增量快照,不用额外配置。但要注意,如果用的是HDFS的路径,尽量别放到根目录下,权限和磁盘配额经常会出问题。
3.4 从0到1完整运行流程
如果你第一次跑这个链路,我建议按下面这个顺序操作,避免反复踩坑:
- 在Kafka里建好topic,确认分区数。我们订单量大,topic设了12个分区,能保证Flink SQL阅读时的并行度至少可以到12。
- 用Flink CDC的datastream方式或者直接先跑一个SQL任务把全量数据先灌到Kafka,这一步称为“全量初始化”。Flink CDC连接器支持全量加增量,它会先做一次快照,再实时监听Binlog。
- 启动Flink SQL任务前,先确认ClickHouse表已经用ReplacingMergeTree引擎建好,并且版本列用Binlog的时间戳,这样即使偶尔重复写入了也能在Merge阶段去重。
- 提交任务,观察Flink UI的指标。重点看Source端的currentFetchEventTimeLag,如果这个值持续走高,说明消费延迟了,需要增加并行度。
- 到ClickHouse里跑一条count确认数据量和延迟。我们验收标准是:从MySQL写入到ClickHouse可见,延迟不超过5秒。
整套跑通之后,运维变得很简单。业务方需要加字段,只需要改Flink SQL的建表语句,然后保存点重启,不用动ClickHouse表结构,不用动Kafka配置。这就是Flink SQL的吸引力:改SQL就够了。
4. 真实项目中高频踩坑与排查实录
4.1 Flink JDBC连接器异常,大多数是配置问题
热搜里“flink的jdbc连接器异常”是高频痛点。我遇到过几次,最常见的原因有三个。
第一个是驱动版本冲突。Flink的JDBC连接器默认自带的MySQL驱动版本很老,如果你在lib目录里又放了一个高版本驱动,两者冲突,启动时报NoClassDefFoundError。解决办法:把Flink lib目录下自带的老版驱动删掉,只保留项目里引入的版本。
第二个是参数名写错。很多人会把JDBC连接器的参数名记混,比如写成:
WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://...', 'table-name' = 'xxx', 'user' = 'root', 'password' = '123456' )这里单值形式写user、password在有新版连接器里可能不生效,应该用'username' = 'root', 'password' = 'xxx'。如果你不确认版本,直接看官方文档对应的连接器参数项,别靠记忆。这个错误最坑的地方是任务能启动,但运行一会儿才报Access denied for user,排查半天才发现是参数名的问题。
第三个是连接数耗尽。JDBC连接器每个并行子任务都会建立连接,如果上游并行度设了20,但是MySQL侧max_connections只有100,其他应用再占一些,就会报Connection is not available, request timed out。解决办法给MySQL加连接数,或者调低source并行度,同步场景尽量用CDC连接器而不是JDBC轮询。
4.2 数据延迟越来越大,先查反压再查并行度
有一次我们的实时大屏数据延迟从5秒涨到15分钟,我打开Flink UI看,发现Source端和Sink端之间有严重的背压。反压的意思是下游处理不过来,上游只能停下来等待,整个管道被堵住了。
排查思路分两步。第一步看是哪个算子反压,通常反压会向上游传播。当时我们发现ClickHouse Sink算子反压最严重,说明瓶颈在写入端。第二步优化写入参数,把sink.batch.size从500调到2000,sink.batch.interval从2秒调到5秒,让更多数据攒成一个批次再写,ClickHouse写入次数减少,压力瞬间下降。反压消除,延迟降到3秒以内。
如果你遇到的情况是Source端持续繁忙但Sink端空闲,那就说明是读取瓶颈,优先增加source并行度。但注意并行度不能超过Kafka的分区数,否则多出来的并行子任务会闲置,没意义。
4.3 窗口结果不输出,先确认watermark有没有推进
新手最容易犯的错是窗口一直不触发。事件时间的滚动窗口必须等到watermark越过窗口结束时间才会触发。如果数据长时间不更新,watermark不推进,窗口就永远悬在那里。
排查方法很简单,在Flink UI的"Watermark"列看当前值。如果一直是-9223372036854775808(Long的最小值),说明上游没有定义watermark或者source里没有把时间字段声明为事件时间。如果watermark确实在涨但窗口还是不出结果,检查一下数据的事件时间字段是不是被当作字符串了,需要先TO_TIMESTAMP()转换再赋值给WATERMARK。
我还遇到过一个很隐蔽的问题,Kafka分区里的数据时间戳严重乱序,比如最早和最晚相差几小时,那么order_time - INTERVAL '5' SECOND的watermark会一直卡在最早那条数据的时间点,导致窗口大面积延迟触发。后来我们把watermark策略改成了按分区允许乱序,并且在SQL任务里加上SCAN.STARTUP.MODE = 'latest-offset'跳过历史脏数据,问题才解决。
4.4 状态无限增长,要定期清理和主动设置TTL
Flink SQL的group聚合如果没有时间范围,状态会无限增长。比如前面那个GROUP BY user_id的累计统计,每个用户ID都会占用状态,用户量涨到几亿,状态也涨到几亿条。这是很多任务状态爆炸的真正原因。
解决办法分两个层面。如果业务允许,聚合加上时间窗口,让状态周期性地随着窗口过期自动清理。如果不允许,那么给状态设置TTL。Flink SQL里可以在建表语句或者环境配置中指定状态TTL:
SET table.exec.state.ttl = 1h;这里的意思每个key的状态如果1小时没更新就自动过期清除。我就吃过亏,有个累计UV统计没设TTL,状态从几GB一路涨到30GB,Checkpoint频繁超时。加上TTL之后,任务稳定运行,状态控制在5GB以内。
4.5 一个花钱买来的教训:并行度不要乱调
刚开始用Flink SQL时,我觉得并行度开得越高越快,于是把source并行度调到Kafka分区数的3倍。结果不仅没变快,反而性能下降了。原因是每个并行子任务都要建立独立的Kafka消费连接和网络连接,连接数越多,协调开销越大。
正确做法是:source并行度等于Kafka分区数;中间算子并行度除非有数据倾斜问题,否则跟source一致;sink并行度看下游性能,ClickHouse一般2-4个并发就能发挥不错的查询性能了。并行度不是越高越好,而是匹配资源的最优状态。上生产前,花半天用小流量实测各个并行度组合下的吞吐和延迟,比上线后拍脑袋调参靠谱得多。
4.6 常见问题速查表
我把自己维护的几个实时链路问题整理成了速查表,团队里新人遇到问题先查表,效率高很多。
| 现象 | 可能原因 | 排查方向 |
|---|---|---|
| 任务启动报ClassNotFound | 连接器Jar包版本冲突或缺失 | lib目录jar排查、maven依赖树 |
| 数据重复写入 | 缺主键去重、重启后状态丢失 | 检查SQL是否有row_number去重、Checkpoint是否开启 |
| 窗口不触发 | watermark未推进、时间字段类型错误 | Flink UI看watermark、检查DDL时间字段 |
| 延迟不断上涨 | 反压、sink写入慢 | UI看反压位置、调batch参数 |
| Checkpoint超时 | 状态过大、HDFS写入慢 | 改RocksDB、开增量、设TTL |
| 连不上下游数据库 | 连接数耗尽、网络不通 | 排查连接池参数、看下游端错误日志 |
4.7 一个隐藏很深的问题:SQL任务重启后offset怎么定
很多团队把Flink SQL任务重启后,Kafka消费位点从earliest开始重放,结果数据重复一大堆。问题在于Checkpoint配置了,但重启时没有用保存点恢复。正确做法是在停止任务前先做一次Savepoint,然后重启时指定--fromSavepoint路径恢复。这个操作如果不会做,那你每次重启都在双写数据。
我在线上推过一个制度:任何Flink SQL任务上线前必须写清楚停止和恢复的操作手册,用保存点恢复是标配。同事们一开始觉得麻烦,直到有人因为没做保存点导致一周的实时数据全部重放、下游数仓被搞乱,大家才意识到这不是可选项,而是必须项。
最后再分享一个小技巧
如果你在测试环境想快速验证一套Flink SQL逻辑,不要直接起集群,用Flink SQL Client本地模式,或者用Flink 1.16之后自带的SQL Gateway,配一个测试Kafka topic,SQL逻辑几分钟就能验证完。我自己的习惯是维护了三个目录:ddl目录放所有表的建表语句,etl目录放所有的insert SQL,udf目录放自定义函数。线上改SQL只需要改对应那一个文件,历史版本清晰可追溯。这习惯帮我解决了很多次“这个SQL是谁改的、为什么改”的排查难题。
Flink SQL的学习曲线不算陡,但真正的深水区在原理理解、参数调优和问题排查。如果你正在做或准备做实时数仓,建议把窗口函数、watermark、状态管理和连接器参数吃透,这四个点搞定,绝大多数场景你都能稳得住。剩下的就是多踩坑、多总结,把每次线上事故变成团队的知识资产,这才是最有价值的成长路径。