年初接手实时数仓项目时,我把自己之前用DataStream API写的一套实时指标任务全部重构了一遍。重构完成后,代码量从两千多行降到四百多行,开发周期从两周缩到三天。让我下定决心做这件事的,正是Flink SQL。这篇不是官方文档复读机,而是我在生产环境里把Flink SQL从选型、建模、同步到ClickHouse,再到排查JDBC连接器异常的一整套实战记录。如果你已经会用DataStream API,但始终没把Flink SQL用起来,或者刚接触大数据实时计算、想找一份能落地参照的案例,这篇文章值得你花二十分钟读完。
1. 先用SQL去写流处理,到底图什么
很多人一听到“Flink SQL”,第一反应是“这不就是把流数据查一下吗,哪有DataStream API灵活”。这个想法不算错,但生产级的灵活并不体现在“所有事情都手写”上。我见过太多团队用DataStream API写了大量窗口聚合、双流关联、状态管理,最后线上出问题,排查一堆时间找不出原因。SQL的引入,本质上是把“你关心什么”和“怎么实现”分开。
1.1 从DataStream API到SQL:阈值在哪
先说清楚边界。如果任务只是简单的读Kafka、过滤、写MySQL,DataStream API一点问题没有,代码也短。但一旦涉及多流关联、滚动窗口、会话窗口、维表补全、迟到数据处理,你要手动处理的细节就会指数级上升:timers怎么注册、状态该存在哪个key、watermark怎么在算子和窗口之间传递、迟到数据能不能被兜住。这些逻辑在SQL里就是几句声明,Flink的查询优化器会负责生成执行计划、调整join顺序、选择合适的state backend。
我自己的判断标准是:任务里只要出现两个join以上,或者有按时间维度的聚合,那用DataStream API手写就划不来了。SQL虽然不能在每个细节上精准控制,但稳定性和可维护性要远高于手写状态代码。实时任务最贵的不是CPU,而是“半夜报警后花两小时查逻辑”。
1.2 Flink SQL能替代什么,替代不了什么
能替代的场景非常明确:
- 实时数据清洗和格式转换:JSON解析、类型转换、异常值过滤。
- 实时指标计算:按分钟、小时进行PV、UV、GMV统计。
- 双流关联:订单流和支付流、订单流和发货流。
- CDC数据同步:把MySQL binlog转成流,再写入其他存储。
- 维表关联:用户维度、商品维度补全,用lookup join搞定。
替代不了的场景也有:
- 需要自定义state结构的复杂风控规则,比如一台设备上五分钟内关联了多个账号,需要用图结构去存。
- 强依赖自定义窗口语义,比如窗口不是按时间闭区间,而是要按业务里的“一单完结”事件触发。
- 对背压和checkpoint有极致控制要求的低延迟场景。
不过这些场景占比很小。我所在的业务线,实时任务里大概80%都在SQL能力范围内。剩下20%才拆出来用DataStream API。
1.3 实时数仓分层里的SQL位置
Flink SQL最常见的落地方式是实时数仓分层。OSS层日志先进Kafka,SQL任务消费Kafka做解析落到Hive或者Iceberg;DWD层用SQL做清洗、维表关联、双流JOIN;DWS层用SQL做分钟级、小时级聚合结果写ClickHouse或Doris。这条链路里,除了最底层的数据接入偶尔要写SourceFunction,剩下全是SQL。
举个例子,我们有个“用户点击行为实时分析”的需求,原链路是:DataStream API读取点击日志 → 做IP解析 → 开三分钟窗口 → 去重 → 聚合成指标。这一套下来要写四百多行Scala。换成SQL之后,建一张点击日志源表、一张IP维表、一张目标结果表,中间加一条CREATE VIEW和一条INSERT INTO就结束了。代码少了,粒度还更清晰,后面同事接手也更容易理解。
2. 搭一套能跑的Flink SQL开发环境
Flink SQL看起来只是写几句SQL,但真要落到自己的项目里,环境搭建和依赖选择是第一个坑。版本选错、依赖缺包、作业提交方式不合适,都会让你在SQL还没跑起来前就消耗掉热情。
2.1 版本不要乱选,先对齐生态
我推荐直接从Flink 1.17或1.18起步,Java环境用8或11都行,但尽量用11。原因很简单:官方连接器和Table API在1.13之后才进入稳定期,1.15之后语法和配置项与现代用法比较接近,再老的版本很多特性缺不说,社区资料也少。
这里列一份我常用的依赖清单,你直接在Maven里配好:
| 依赖 | 作用 |
|---|---|
| flink-table-api-java | Java版的Table API和SQL API |
| flink-table-planner-loader | 内置查询优化器,作业运行时需要 |
| flink-table-runtime | SQL执行时的运行时组件 |
| flink-connector-kafka | Kafka Source/Sink连接器 |
| flink-connector-jdbc | JDBC Sink连接器 |
| mysql-connector-java | MySQL驱动 |
| clickhouse-jdbc | ClickHouse驱动 |
有一件事特别提醒:Flink的flink-table-planner和flink-table-runtime如果单独引入,版本必须和Flink主版本完全一致。我遇到过很多次本地能跑、提交集群就报NoClassDefFoundError,十有八九是planner-loader和表的运行时依赖没对齐。
2.2 代码写SQL还是SQL客户端写SQL
我本人习惯先用Flink SQL客户端做语法验证,再把SQL以字符串形式固化到Java代码里。这样既能快速调试,又能把最终SQL纳入版本控制。
如果只是临时分析,直接启动sql-client.sh,写CREATE TABLE和INSERT INTO试试就行。但生产任务建议都写成Java代码,因为你多半要配合自定义UDF、运行时参数和状态恢复机制。用Java代码写SQL有一个好处:可以用StreamTableEnvironment.executeSql()一条一条提交,逻辑清晰,也方便做多级ETL的编排。
2.3 一个最小可跑的Java任务
下面这个是我每次新建Flink SQL项目时的起点,读Kafka、打印结果,验证环境和依赖没问题后再往里面加业务:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); EnvironmentSettings settings = EnvironmentSettings.inStreamingMode(); StreamTableEnvironment tEnv = StreamTableEnvironment.create(env, settings); tEnv.executeSql( "CREATE TABLE source_event (\n" + " user_id BIGINT,\n" + " event_type STRING,\n" + " event_time BIGINT,\n" + " ts AS TO_TIMESTAMP_LTZ(event_time, 3),\n" + " WATERMARK FOR ts AS ts - INTERVAL '5' SECOND\n" + ") WITH (\n" + " 'connector' = 'kafka',\n" + " 'topic' = 'user-event',\n" + " 'properties.bootstrap.servers' = 'localhost:9092',\n" + " 'properties.group.id' = 'flink-sql-demo',\n" + " 'scan.startup.mode' = 'latest-offset',\n" + " 'format' = 'json'\n" + ")" ); tEnv.executeSql( "CREATE TABLE print_sink (\n" + " user_id BIGINT,\n" + " event_type STRING,\n" + " cnt BIGINT\n" + ") WITH (\n" + " 'connector' = 'print'\n" + ")" ); tEnv.executeSql( "INSERT INTO print_sink\n" + "SELECT user_id, event_type, COUNT(*)\n" + "FROM source_event\n" + "GROUP BY user_id, event_type" ); env.execute("flink-sql-demo");这段代码跑起来后,你会看到控制台每秒输出聚合结果。等这一步通了,再往上加真实业务就不容易产生“环境问题”。
3. 把Kafka数据接进来:动态表与连续查询
第一次真正理解Flink SQL和普通数据库SQL的区别,是在排查一个聚合任务为什么一直输出不完整时。普通SQL查的是静态的数据集,查完就结束了;Flink SQL查的是一张永远在增长的“动态表”,查询也永远不会停止,除非作业被取消。这个本质差异,决定了建表、写查询、调状态的整个思路。
3.1 流和表到底怎么映射
Kafka里一条一条消息,本质上是一个无限增长的流。Flink SQL把它看作动态表:插入一条消息,相当于往表里加一行;删除或者更新一条,也对应表的变更。动态表有三种模式:
- Append-only:只插入,不更新。适合日志、事件流。
- Retract:既可以插入也可以删除。适合
COUNT、SUM这类会产生回撤的聚合。 - Upsert:按主键更新。适合维表、订单状态表。
这点非常关键。比如你执行SELECT user_id, COUNT(*) FROM source_event GROUP BY user_id,因为输入流不断有新数据,同一个user_id的COUNT(*)会不断变大,Flink SQL就需要通过Retract消息把上一轮的旧结果作废,再输出新结果。如果你用print sink调试,会看到先输出一个“-1”,再输出“+2”,这代表撤回和新增。
3.2 建表语句里的格式与时间字段处理
Kafka里的数据一般是JSON字符串。建表时除了指定connector和topic,最容易被忽略的是时间字段。
比如原始日志里事件时间是毫秒级时间戳:
CREATE TABLE user_event ( user_id BIGINT, event_type STRING, event_time BIGINT, ts AS TO_TIMESTAMP_LTZ(event_time, 3), WATERMARK FOR ts AS ts - INTERVAL '5' SECOND ) WITH (...)这里ts AS TO_TIMESTAMP_LTZ(event_time, 3)是把BIGINT类型的事件时间转成TIMESTAMP_LTZ,WATERMARK则告诉Flink最多容忍5秒乱序。没有这一步,窗口聚合基本跑不出准确结果。
还有一个小提示:Kafka JSON格式解析时,尽量加'json.ignore-parse-errors' = 'true'。生产环境里上游同事偶尔会发一条脏数据,如果不开忽略解析错误,整个作业会因为一条坏消息直接失败,然后反复重启。
3.3 连续查询不是“跑一次就结束”
普通SQL对静态表查一次,返回结果就完了。Flink SQL的连续查询会一直运行,只要输入表有数据进来,就会重新计算并更新结果。这就意味着状态会持续累积。
拿GROUP BY为例,Flink必须把每个key的当前存量一直保存在状态里。如果条件不合适,比如你按用户ID分组,用户量上亿,状态就会非常大。好的一点是Flink有状态生存时间配置,你可以在建表或者作业配置里设置table.exec.state.ttl,超过时间没更新的key自动清理,防止状态无限增长。
我在生产环境常用的设置是:
SET 'table.exec.state.ttl' = '2 h';这里的2小时需要结合业务合理设置。如果窗口聚合是1小时周期,那2小时TTL足够;如果业务需要精确去重且去重周期为7天,TTL就必须不小于7天。
4. 生产最常用的三类Flink SQL场景:过滤、聚合、双流JOIN
Flink SQL真正产生价值的地方,还是那些在实时数仓里高频出现的场景。我总结下来就三类:清洗过滤、窗口聚合、流与流的关联。把这三类吃透,你就能覆盖绝大多数实时指标需求。
4.1 清洗与字段治理
实时清洗不是只做WHERE这么简单。你经常要处理类型不一致、字段缺失、枚举值混乱,甚至同一个日志里有新旧两版格式并存。
一种常见写法是用CASE WHEN做标准化:
INSERT INTO dwd_event SELECT user_id, LOWER(COALESCE(event_type, 'unknown')) AS event_type, TO_TIMESTAMP_LTZ(event_time, 3) AS event_ts, CASE WHEN platform = 'ios' OR platform = 'IOS' THEN 'iOS' WHEN platform = 'android' THEN 'Android' ELSE 'other' END AS platform FROM source_event WHERE user_id IS NOT NULL AND event_time > 0;注意COALESCE只能把NULL转成默认值,如果是空字符串,还需要单独处理。另外,如果源Kafka的JSON字段和表字段类型对不上,Flink SQL默认会直接报错,这时候可以在WITH里加'format'='json'配合'json.fail-on-missing-field'='false',允许缺失字段,能让你省下很多折腾解析的时间。
4.2 窗口聚合的语法陷阱
窗口聚合是实时指标最常见的操作。Flink SQL里最常用的滚动窗口是TUMBLE:
SELECT TUMBLE_ROWTIME(ts, INTERVAL '5' MINUTE) AS window_time, event_type, COUNT(*) AS cnt FROM source_event GROUP BY TUMBLE(ts, INTERVAL '5' MINUTE), event_type;这里有个很容易踩的语法坑:GROUP BY里必须同时带上TUMBLE(...),否则Flink不知道数据要按哪个时间列开窗。窗口函数返回的window_time有特殊的类型,写TUMBLE_ROWTIME或TUMBLE_END都可以,但别把两者混用,否则结果表里的时间字段含义会不一致。
事件时间的乱序处理也值得单独说。如果你在源表上定义了5秒的watermark,那么在窗口末尾5秒内到达的数据,会先触发窗口计算,之后又来一批5秒前的事件,Flink不会重新更新已关闭的窗口结果。想要处理更晚的数据,需要额外用侧输出流或ALLOWED LATENESS语法,但生产环境里为了简单,我一般把watermark的容忍范围设成日志到达延迟的90分位值。
4.3 双流JOIN怎么避免状态爆炸
实时双流JOIN是很多人头痛的问题。两个Kafka流直接JOIN,如果没有任何时间边界,Flink就必须把所有历史数据都放在状态里,数据量一大直接OOM。所以99%的生产任务都会用interval join,只关联某个时间窗口内的数据:
SELECT o.order_id, o.user_id, p.pay_amount FROM orders o JOIN payments p ON o.order_id = p.order_id AND p.pay_time BETWEEN o.order_time AND o.order_time + INTERVAL '10' MINUTE;这个BETWEEN ... AND ...是核心,它告诉Flink只需要为每个订单保留10分钟内的支付记录状态,超过之后自动清理。很多实时数仓项目里的“近N分钟活跃用户”“实时转化漏斗”,底层其实都是这种interval join。
有一点要提醒:如果两个流的order_id分布极度不均匀,比如某个头部用户占了大量订单,JOIN算子仍然会倾斜。处理办法是给订单ID加盐拆散计算,或者把热点key单独分流,这个后面会详细说。
4.4 用UDF补齐SQL覆盖不到的逻辑
Flink SQL内建函数很多,但不是一切都能覆盖。比如我们业务里要根据IP解析城市、根据UA解析设备,这些就得用UDF。做法很简单,把函数注册进Table Environment:
tEnv.createTemporaryFunction("ip_to_city", IpToCityUdf.class);然后在SQL里直接用SELECT ip_to_city(ip) FROM source_event。UDF能极大拓展SQL的边界,但注意不要在UDF里做一些阻塞式的远程RPC调用,否则会把整个TaskManager线程卡住。如果需要查维表,优先用lookup join而不是UDF里调接口。
5. 从MySQL同步到ClickHouse:一套Sink链路实战
实时数仓里最有代表性的落地案例之一,就是“MySQL数据实时同步到ClickHouse”。我们的订单表在MySQL,运营看板在ClickHouse,中间需要毫秒级延迟地搬数据。这套链路用Flink SQL写起来非常清爽。
5.1 同步需求拆解
订单表有三张:订单主表、订单明细表、订单状态流水表。ClickHouse里对应的表结构不一样,需要按维度宽表建模。数据同步不能只是简单复制,还要把状态流水合并进去,按主键更新最新状态。
这就涉及三类操作:读取MySQL变更、按订单ID关联、写入ClickHouse并更新。
5.2 读取MySQL的两条路线
路线一是直接用Flink的JDBC连接器轮询查询:
CREATE TABLE mysql_orders ( id BIGINT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10, 2), status STRING, update_time TIMESTAMP(3) ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://localhost:3306/business', 'table-name' = 'orders', 'username' = 'root', 'password' = '123456' );这能做,但JDBC source是扫描查询,不是binlog级实时变更捕获,会有分钟级延迟,而且每次扫描全表成本很高。
路线二是用Flink CDC:
CREATE TABLE mysql_orders ( id BIGINT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10, 2), status STRING, update_time TIMESTAMP(3) ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'localhost', 'port' = '3306', 'username' = 'flinkuser', 'password' = '123456', 'database-name' = 'business', 'table-name' = 'orders' );CDC方式会直接消费binlog,数据是真正的实时变更,还能识别数据的新增、更新和删除。我们最终选的是CDC。这里提醒一下,mysql-cdc是独立连接器,需要在pom里额外引入flink-connector-mysql-cdc,不是Flink官方自带。
5.3 Sink到ClickHouse:用JDBC还是ClickHouse连接器
Flink官方其实没有提供ClickHouse连接器,社区里有一个flink-clickhouse-connector,但成熟度参差不齐。我评估之后选了最稳的JDBC连接器,加上ClickHouse官方JDBC驱动。
JDBC Sink的建表语句长这样:
CREATE TABLE clickhouse_order_sink ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), status STRING, update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:clickhouse://localhost:8123/default', 'table-name' = 'order_total', 'username' = 'default', 'password' = '', 'sink.buffer-flush.max-rows' = '1000', 'sink.buffer-flush.interval' = '2s', 'sink.buffer-flush.max-retries' = '3' );重点参数是sink.buffer-flush.max-rows和sink.buffer-flush.interval。ClickHouse对批量写入的吞吐要求比较高,如果每一条来一条就写一条,性能会非常差。我们设置成攒到1000条或2秒刷一次,单并发写入速度提升非常明显。
ClickHouse实现更新,我使用的是ReplacingMergeTree表引擎加复合主键,这样即使Flink发送了重复数据,ClickHouse也能通过后台合并去重。如果你想用JDBC sink实现真正的精确更新,需要自己处理幂等逻辑,JDBC连接器本身不保证exactly-once。
5.4 完整同步SQL
整个同步任务一共分三步:定义CDC源表、定义ClickHouse结果表、执行插入:
INSERT INTO clickhouse_order_sink SELECT id, user_id, amount, status, update_time FROM mysql_orders WHERE status <> 'cancel';就这么简单。任务跑起来后,MySQL里改一行订单数据,ClickHouse里大概一秒内就能查到更新。如果还要把订单明细和订单主表关联成宽表,再建一个中间视图,把两个源表做interval join之后统一写入ClickHouse即可。
5.5 同步延迟和精确性
这套链路做不到的是真正意义上的“exactly-once写入ClickHouse”。原因有两个:一是JDBC sink对ClickHouse的支持没有内置两阶段提交,二是ReplacingMergeTree本身是最终一致。我们通过两个办法把误差缩小到可接受范围:
- 写入前在SQL里对主键做
DISTINCT降重。 - ClickHouse表设置
ReplacingMergeTree并指定版本字段,让重复数据后台自动合并。
如果你需要严格不重复,建议在ClickHouse写入端再挂一个去重引擎或用外部字典。实时数仓领域里,能做到最终一致已经很好了,别为了绝对一致把自己困在无法上线的死胡同里。
6. 踩坑实录:JDBC连接器异常和并行度陷阱
看了不少文章,大家都爱讲“Flink SQL多方便”,但实际生产环境里,连接器异常和并行度设置才是真正让人头皮发麻的地方。我把自己线上遇到过的几个高频问题整理出来,希望能帮你们少走弯路。
6.1 ClassNotFoundException和NoClassDefFoundError
SQL作业本地可以跑,一上集群就报ClassNotFoundException: com.mysql.cj.jdbc.Driver,这种是最常见也是最容易解决的。原因是你的Flink集群/lib目录下没有MySQL驱动包,只有本地IDE里带了依赖。
解决方法是把驱动jar包直接放到TaskManager的lib目录,或者提交作业时用-j /path/to/mysql-connector-java.jar显式指定。注意同一个连接器对应的driver版本要一致,我曾经在集群lib里放了MySQL 8.x驱动,但代码里引的是5.x,结果报Communications link failure,定位了好几个小时。
ClickHouse因为同样用JDBC,也要保证clickhouse-jdbc.jar在每台TaskManager上可见。如果用了Flink on YARN部署,最容易的处理方式是放在HDFS上,然后提交时通过配置把jar分发到所有TaskManager。
6.2 Connection is not available,request timed out
这个错误很吓人,我第一反应是ClickHouse数据库挂了,结果排查下来发现是连接池被耗尽。Flink SQL的JDBC Sink默认会为每个并行子任务维护一个连接池,每个连接池的默认大小有限。如果Sink并行度是32,ClickHouse能承受的连接数是50,那么所有TaskManager同时写的时候就会有几批任务拿不到连接。
解决办法有两个方向:
- 调低Sink并行度:比如从32降到8,减少同时打开的连接数。
- 调整用户侧连接池和ClickHouse
max_connections参数。
Flink SQL里不能直接给单表设置并行度,但你可以在整个作业设置execution.parallelism,或者给Sink表单独包装成一条INSERT INTO,然后通过PARALLELISM提示来控制。比如:
INSERT INTO clickhouse_order_sink /*+ OPTIONS('sink.parallelism'='8') */ SELECT ...;如果发现仍然不够,再配合调大sink.buffer-flush.interval,让写入更批量、更平缓。
6.3 并行度和Kafka分区数不匹配导致的数据空转
Kafka topic有12个分区,Flink SQL任务并行度设成了24。Kafka Source最多只有12个并行子任务读数据,剩下12个子任务闲置。这个问题不致命,但会出现一种“看着有24个并发,实际吞吐只有12个分区的能力”的错觉。反过来如果并行度小于分区数,则会同时出现数据倾斜和状态膨胀。
一般我的做法是:Kafka Source并行度等于min(partition数, max(4, 分区数)),聚合和Sink并行度再独立设置。遇到热点key就用SQL的DISTRIBUTE BY或GROUP BY次数扩展,分两步做聚合。
6.4 时区问题导致的日期漂移
这个绝对算隐蔽坑。Flink SQL默认的本地时区是UTC,如果你的Kafka数据里时间戳是“2025-01-01 12:00:00”这样的本地时间,直接转成TIMESTAMP并在窗口聚合,你会发现窗口边界全都偏了8小时。
解决办法是在作业启动参数里明确设置时区:
SET 'table.local-time-zone' = 'Asia/Shanghai';或者在Java代码里创建Environment时配置Configuration。时区问题不解决,所有按天、按小时聚合的指标都会对不上。这是我见过最容易被忽视、但影响最严重的一个“配置参数”。
6.5 Spring Boot整合Flink时的类加载器坑
因为业务平台是Java后端,很多团队想把Flink SQL任务直接嵌到Spring Boot进程里跑,不用独立提交到集群。这在技术上可行,但坑非常多。
最典型的问题是类加载器冲突:Spring Boot用FatJar方式加载类,Flink的Table环境中很多类是懒加载,一旦StreamTableEnvironment在Spring Bean初始化时被创建,后续的动态类加载很可能找不到Flink的planner类,导致TableException或ClassCastException。
我的经验是,如果只是要管理SQL作业,不要直接在Spring Boot的Tomcat类加载环境里去跑Flink集群。可以把Spring Boot应用当成一个“SQL任务提交器”,只负责把SQL文件内容组装成提交命令,真正执行还是交给独立的Flink集群。如果必须在进程内运行,建议单独写一个FlinkSqlRunner类,在子线程里用隔离类加载器启动,不要用Spring注解自动注入StreamTableEnvironment。
7. 上线前必看的State与性能调优手段
Flink SQL任务在开发环境跑通很容易,上线之后能不能扛住高峰期数据量,才是考验。以下这些调优点,都是我根据线上实例反复调参总结出来的。
7.1 Checkpoint配置不能省
状态是Flink SQL任务的心脏,checkpoint是保证状态不丢的前提。我在生产环境里至少会开:
execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE state.backend: rocksdb state.checkpoints.dir: hdfs:///flink/checkpoints间隔不建议太短,比如10秒一次,在高吞吐场景下会导致频繁的barrier对齐,反而降低吞吐。60秒是我常用的区间。如果任务状态很大,把RocksDB启成增量模式能显著减少checkpoint时长:
state.backend.rocksdb.incremental: true7.2 资源并行度估算公式
我给团队定的一个粗估公式:Source并行度约等于Kafka分区数;聚合节点并行度等于QPS * 窗口内key的复杂度 / 单核心吞吐,但更简单的做法是先按Kafka分区数的1.5倍起,再观察背压;Sink并行度则依据下游写入阈值来定。
例如我们的订单同步任务,Kafka分区是24,MySQL CDC任务并行度24够用。ClickHouse Sink如果并行度也是24,每秒能处理超过1万条。如果ClickHouse出现写入瓶颈,不是一味加并行度,而是先调大批量刷新参数。
7.3 MiniBatch和LocalAggregate是聚合任务的救星
默认Flink SQL聚合是一条一条计算,每次来一条数据都要访问状态。数据量大时,状态访问会成为瓶颈。开启MiniBatch后,Flink会在聚合算子内先攒一批数据,再统一访问状态,减少重复读写。
table.exec.mini-batch.enabled: true table.exec.mini-batch.size: 20000 table.exec.mini-batch.latency: 2s同时打开本地聚合:
table.optimizer.agg-phase-strategy: TWO_PHASE这个组合能极大缓解高基数聚合的性能问题。我们在计算用户行为指标时,开启前单作业CPU达到70%,开启后降到40%,延迟也稳定多了。
7.4 数据倾斜的通用解法
Flink SQL聚合最怕某一个key的数据量特别大。比如实时大屏按省份统计点击量,北京一个省的数据可能占到40%。这时候不管并行度开到多少,所有数据都压到同一个key的聚合算子,其他key的空闲。
一种常用的SQL层解法是“两阶段无Key聚合”:先给数据带一个随机后缀,partition到一个临时的局部聚合,再把后缀去掉做全局聚合。Flink SQL里可以这样写:
INSERT INTO result SELECT province, SUM(cnt) AS total FROM ( SELECT province, CONCAT(province, '_', MOD(CAST(ROWNUM AS INT), 10)) AS random_key, COUNT(*) AS cnt FROM source_table GROUP BY province, CONCAT(province, '_', MOD(CAST(ROWNUM AS INT), 10)) ) GROUP BY province;这个做法牺牲了一点精确度,但对大部分指标误差很小。如果必须完全准确,就得在DataStream层自己写keyby的加盐逻辑。
7.5 监控指标比日志更重要
上线不是结束,而是另一种开始。Flink SQL作业我一般盯四个指标:currentInputWatermark判断数据源是否断流、numRecordsInPerSecond判断吞吐是否达标、numBufferedRecords判断Sink是否反压、lastCheckpointSize判断状态增长是否异常。这四个里任何一个出现明显异常,我都会第一时间看SQL的执行计划,而不是先去翻日志。
举个例子,有一次lastCheckpointSize突然涨了3倍,日志里一切正常。后来检查发现是源表某个上游业务字段脏数据太多,导致GROUP BY的key基数飙升。要不是提前盯住了checkpoint大小,等到状态过大再恢复作业时,整条链路可能已经停了一个小时。
最后再分享一个小技巧
这套Flink SQL链路在我们线上稳定跑了近半年,最大并行度到48,每天处理亿级事件。踩了这么多坑以后,我最大的体会是:Flink SQL不是“低配版流处理”,而是一种边界明显的工程选择。把合适任务交给SQL,把复杂状态剥离给自定义算子,你会省下大量时间。最后补一个比JDBC连接器更隐蔽的坑:如果线上任务用了GROUP BY做大窗口聚合,并且状态一直增长,别只想着加资源,先去查状态里的key是不是因为某个字段值本身就一直在变。一个字段写错,状态增长的速度会远远超出你的预期。