Flink状态作业调优实践指南FlinkSQL作业篇
Apache Flink SQL 作业大状态导致反压的调优原理与方法
以下文章来源于 Apache Flink,作者 Flink 团队@阿里云

Apache Flink 是 Apache Flink 中文社区唯一官微,由 Flink PMC 维护。
5.1 运行原理:状态算子的产生
5.1.1 基于优化器推导产生的状态算子
| 状态算子 | 状态清理机制 |
|---|---|
| ChangelogNormalize | 生命周期 TTL |
| SinkUpsertMaterializer | |
| LookupJoin(*) |
(1)ChangelogNormalize
ChangelogNormalize 作为一个状态算子,旨在对涉及主键语义的数据变更日志进行标准化处理。通过这一算子,可以有效地整合和优化数据变更记录,确保数据的一致性和准确性。该状态算子会在以下两种场景出现:
- 使用了带有主键的 upsert 源表
- 用户显式设置
table.exec.source.cdc-events-duplicate= 'true'
(2)SinkUpsertMaterializer
SinkUpsertMaterializer 是一种状态算子,专门用于处理具有主键定义的结果表,并确保数据的物化操作符合 upsert 语义。在数据流更新过程中,如果无法保证 upsert 的特定要求,即按照主键进行更新时保持数据的唯一性和有序性,优化器会自动引入此算子。
(3)LookupJoin
在处理 LookupJoin 操作时,若用户主动配置了系统优化选项 table.optimizer.non-deterministic-update.strategy 为 'TRY_RESOLVE',且优化器识别到潜在的非确定性更新问题,则系统会尝试采取特殊措施以解决这一问题。
5.1.2 基于 SQL 操作产生的状态算子
基于 SQL 操作产生的状态算子,按状态清理机制可以分为 TTL 过期和依赖 watermark 推进两类。
| 状态算子 | 如何产生 | 状态清理机制 |
|---|---|---|
| Deduplicate | 使用 row_number 语句,order by 的字段必须为 time attribute (event time 或 processing time ),且只取第一条 | TTL |
| RegularJoin | 使用 join 语句,等值条件里不包含 time attribute 字段 | |
| GroupAggregate | 使用 group by 语句进行分组聚合,如 sum/count/min/max/first_value/last_value,或使用 distinct 关键字 | |
| GlobalGroupAggregate | 分组聚合开启 local-global 优化 | |
| IncrementalGroupAggregate | 当存在两层分组聚合操作并开启两阶段优化时,内层聚合对应的状态算子 GlobalGroupAggregate 和外层聚合对应的 LocalGroupAggregate 被合并成一个 IncrementalGroupAggregate | |
| Rank | 使用 row_number 语句,order by 的字段必须为非 time attribute 字段 | |
| GlobalRank | 使用 row_number 语句,order by 的字段必须为非 time attribute 字段,并开启 local-global 优化 | |
| IntervalJoin | 使用 join 语句,等值条件里包含时间属性(time attribute)字段,可以是事件时间(event time)也可以是处理时间(processing time) | watermark |
| TemporalJoin | 使用基于事件时间(event time)的 inner 或 left join 语句 | |
| WindowDeduplicate | 基于 Window TVF 的去重操作 | |
| WindowAggregate | 基于 Window TVF 聚合 | |
| GlobalWindowAggregate | 基于 Window TVF 聚合+开启两阶段优化 | |
| WindowJoin | 基于 Window TVF 的关联 | |
| WindowRank | 基于 Window TVF 的排序 | |
| GroupWindowAggregate | 基于 legacy 语法的 Window 聚合 |
5.2 问题诊断方法
同上节中的诊断方法:Flink 大状态作业调优实践指南:Datastream 作业篇
5.3 调优方法
5.3.1 主动避免生成不必要的状态算子
基于 SQL 操作的状态计算一般很难避免,这里主要针对优化器自动推导的算子进行讨论。
(1)ChangelogNormalize
在使用 upsert source 进行数据处理时,我们需注意其 ChangelogNormalize 这种状态节点的生成。通常情况下,除了事件时间的时态关联(event time temporal join)之外,其他 upsert source 应用场景都会产生该状态节点。
(2)SinkUpsertMaterializer
在 table.exec.sink.upsert-materialize 配置项中,AUTO 作为其预设选项,表明系统会自动判断数据的一致性,尤其是在变更日志(changelog)出现无序的情况下。
5.3.2 减少状态访问频次:开启 mini-batch
在对延时要求不高(比如分钟级别的更新)的场景下,开启 mini-batch 攒批优化将会减少 state 的访问和更新频率,提升吞吐。
5.3.3 减少状态大小:设置合理生命周期
在优化计算系统时,关键在于精简状态数据以提高性能。通过减少不必要的状态信息,我们可以显著提升状态访问的速度。TTL(Time-to-Live)策略在此过程中扮演着重要角色,它通过设定数据的存活时间来控制状态数据的规模。
5.3.4 减少状态大小:命中更优的执行计划
在生成执行计划时,优化器会结合输入 SQL 和配置选择相应的 state 实现。
(1)利用主键优化双流连接
- 当连接键(join key)包含主键时,系统采用 ValueState
进行数据存储,这样可以为每个连接键仅保留一条最新记录,实现存储空间的最大化节省。 - 如果连接操作使用了非主键字段,即使已定义主键,系统会使用 MapState<RowData, RowData> 进行存储,以便为每个连接键保存来自源表的、基于主键的最新记录。
- 在未定义主键的情况下,系统将使用 MapState<RowData, Integer> 存储数据,记录每个连接键对应的整行数据及其出现次数。
(2)优化 append_only 流去重操作
使用 row_number 函数替代 first_value 或 last_value 函数进行去重,可以更有效地保留首次或最新出现的记录。
(3)提升聚合查询性能
在进行多维度统计,如计算全网 UV、手机客户端 UV、PC 端 UV 等时,推荐使用 AGG WITH FILTER 语法替代传统的 CASE WHEN 语法。
5.3.5 减少状态大小:调整多流 join 顺序缓解 state 放大
Flink 在处理数据流时,采用了二进制哈希连接(binary hash join)的方式。在示例中,A 与 B 的连接结果会导致数据存储的冗余,这种冗余程度与连接操作的频率成正比。
5.3.6 尽可能减少读盘
参考“Flink 大状态作业调优实践指南:Datastream 作业篇”中的调优方法。

end
