改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

Flink状态作业调优实践指南FlinkSQL作业篇

ckckck2025年1月10日32 浏览

Apache Flink SQL 作业大状态导致反压的调优原理与方法

以下文章来源于 Apache Flink,作者 Flink 团队@阿里云

Apache 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 作业篇”中的调优方法。

Flink

end