
如果你也在维护一个 24 小时不落地的实时作业那你一定遇到过这种情况任务刚上线时一切正常跑了几天之后Web UI 上的 checkpoint 开始越涨越大反压报警一条接一条。我上次排查一个“简单去重”SQL 时第一反应是把并发调大、把堆内存翻倍结果状态还是照涨不误。最后救我的不是内存是 FlinkSQL 的EXPLAIN PLAN——它非常直观地告诉我这条 SQL 在底层被翻译成了哪些算子、哪个算子攥着状态不放手。这篇 Day 21 就把我看执行计划的方法和“状态算子识别”的经验一次说清楚。这类问题几乎每个用 Flink SQL 做实时计算的人都会撞上。尤其是刚接触状态后端、或者是把 SQL 当成“黑盒”来用的同学通常只在作业出问题时去调state.backend配置却很少回过头看 SQL 到底被优化成了什么。实际上SQL 本身只是入口真正决定状态规模的是执行计划里的一个个StreamExec*算子。想控制状态先得搞清楚执行计划在说什么。1. 为什么要盯执行计划一次真实的状态膨胀排查1.1 事情经过从“加内存”到“查计划”大概去年这个时候我接到一个实时 UV 指标作业。业务侧要求按用户维度去重统计当天访问量SQL 写得也不复杂大概是COUNT(DISTINCT user_id)加一个小时级别的窗口。作业刚上线几天非常稳定延迟几十毫秒。可到了第五天Flink Web UI 里 checkpoint 大小从不到 1GB 涨到了 6GB下游 Kafka 开始积压反压分数到了 0.8。当时团队里第一反应就是状态撑不住了加内存。我把 TaskManager 堆内存从 4GB 调到 8GB并行度从 3 提到 6重启之后确实撑了半天但很快又涨回去了。后来我实在没辙才想起来把生产 SQL 原封不动导出来在 SQL Client 里跑了一次EXPLAIN又去 Web UI 里翻了算子详情问题一下就清楚了SQL 里有三个地方被优化器改写成了带状态的算子而且状态 TTL 使用的是默认配置等于不清理。加内存只是把“病死时间”往后拖了一点根本没碰到病根。1.2 EXPLAIN PLAN 能回答的三类问题很多同学觉得EXPLAIN只是面试题或者文档概念实际排查时宁可去翻日志也不愿意看它。但在我看来执行计划是少数能回答下面三个问题的直接入口这条 SQL 在运行时到底是哪些算子在工作而不是你脑子里想象的“几个步骤”哪些算子是有状态算子哪些只是纯计算节点数据是怎么在算子之间分发的是hash还是rebalance这决定了状态会不会“长歪”。从那次之后我养成了一个习惯任何一条会长期运行的 Flink SQL上线前必须把EXPLAIN输出存档一份。它不复杂也不需要每行都读懂但至少要能回答“这条 SQL 会不会产生状态、产生在哪里”这个问题。2. EXPLAIN() 的三层输出从语法树到物理算子2.1 三层输出分别是什么在 Flink SQL 里执行EXPLAIN SELECT ...或者是在代码里调用tableEnv.explainSql()你一般会看到三段内容Abstract Syntax Tree、Optimized Logical Plan、Physical Execution Plan。Abstract Syntax Tree是 SQL 文本被解析后的语法树最贴近你写的 SQL 结构基本就是投影、过滤、表扫描的简单罗列。它解决的是“SQL 写没写对”的问题对状态排查用处不大。Optimized Logical Plan是经过列裁剪、谓词下推、常量折叠、聚合优化之后的逻辑计划。这里已经能看到关键的GroupAggregate、Join、Rank等节点也能看到数据分布方式Exchange(distribution[hash[...]])。它解决的是“优化器把我的 SQL 改成了什么”的问题。Physical Execution Plan是最终物化出来的运行层算子名字通常长这样StreamExecGroupAggregate、StreamExecJoin、StreamExecRank。这一层才是和状态直接相关的地方因为它对应到真实运行时节点。我们在 Web UI 里看到的执行图基本就是从物理计划生成的。我拿一条很简单的分组聚合作例子SELECT user_id, COUNT(order_id) AS cnt FROM orders GROUP BY user_id;它的EXPLAIN大致会输出 Abstract Syntax Tree LogicalProject(user_id[$1], cnt[$2]) - LogicalAggregate(group[{1}], cnt[COUNT($0)]) - LogicalProject(order_id[$0], user_id[$1]) - LogicalTableScan(table[[default_catalog, default_database, orders]]) Optimized Logical Plan GroupAggregate(groupBy[user_id], select[user_id, COUNT(order_id) AS cnt]) - Exchange(distribution[hash[user_id]]) - TableSourceScan(table[[default_catalog, default_database, orders]], fields[order_id, user_id, amount, ts]) Physical Execution Plan StreamExecGroupAggregate(groupBy[user_id], select[user_id, COUNT(order_id) AS cnt]) - StreamExecExchange(distribution[hash[user_id]]) - StreamExecScan(table[[default_catalog, default_database, orders]], fields[order_id, user_id, amount, ts])注意这里的StreamExecExchange(distribution[hash[user_id]])它对应的是 keyBy本身不是状态算子但它是后续状态算子的“前站”。任何 keyed state 都要求数据按 key 分区所以看到hash分布后面紧跟着聚合、join、rank 类算子基本可以断定下一个算子要产生状态了。2.2 怎么读一个物理算子的描述很多人第一次看StreamExecGroupAggregate这种名字就被劝退了其实拆开很直白StreamExec表示流执行节点GroupAggregate表示分组聚合括号里的groupBy[user_id]告诉你它按哪个字段分组。更详细的状态信息比如状态 TTL、状态前后端类型、是否开启增量检查点通常不会全部塞在EXPLAIN的文本输出里。尤其在新版本 Flink 中想看到完整信息有两种方式用EXPLAIN的 JSON 模式把物理计划以 JSON 形式导出来里面会有每个算子的 uid、chaining strategy 等直接把作业提交到集群去 Flink Web UI 看对应算子的详细信息页。我自己的习惯是先用文本模式的EXPLAIN快速判断“有没有状态、状态在哪”需要精细调参时再去看 Web UI 里的算子详情。不要一开始就把自己埋进 JSON 里信息太多反而看不清。3. 会产生状态的六类 SQL执行计划里的“状态指纹”这一节是重头戏。下面六类 SQL 是我在真实生产环境里逐一验证过的它们在执行计划里的算子名、状态结构、清理语义都不太一样。我总结成“状态指纹”方便大家对着自己的执行计划做匹配。3.1 分组聚合状态里存的是“还没算完的部分聚合结果”标准 SQL 里最常见的GROUP BY对应物理算子StreamExecGroupAggregate。它的状态以分组 key 为键value 是聚合中间结果。COUNT、SUM、AVG这种简单聚合的状态很小每个 key 几十个字节但如果是LISTAGG、COLLECT这种收集类聚合value 会随着数据量线性增长。我之前排过一个慢 SQLSQL 文本就是普通的GROUP BY但里面有一个LISTAGG把用户所有行为标签拼成一串字符串。状态里每个 key 都挂着一个不断变长的字符串跑了几个小时就能吃掉几个 GB。看执行计划时如果发现聚合节点上带着LISTAGG或COLLECT基本就要警惕状态膨胀。另一个关键是如果分组 key 本身是高基数维度比如按user_id分组那就意味着状态里的 key 条数会一直增加直到等于用户总量。这种场景下状态 TTL 就非常重要否则 key 永远不会删除。执行计划里的典型特征GroupAggregate(groupBy[user_id], select[user_id, LISTAGG(tag) AS tags]) - Exchange(distribution[hash[user_id]])看到groupBy后面跟的是长尾高基数 key再叠加收集类聚合这就是一个典型的“状态会失控”的信号。3.2 窗口聚合每个窗口一份缓冲与 watermark 紧密相关窗口聚合在 SQL 里通常是这样写的SELECT TUMBLE_START(ts, INTERVAL 10 MINUTE) AS win_start, user_id, COUNT(*) FROM orders GROUP BY TUMBLE(ts, INTERVAL 10 MINUTE), user_id;对应物理算子一般是StreamExecWindowAggregate。它的状态结构和普通分组聚合不同除了分组 key每个窗口也有自己的一份聚合缓冲。窗口结束并触发计算之后这部分状态会被清理前提是 watermark 正常推进。实际生产里最常出问题的就是 watermark 不推进。比如上游 Kafka 某个 partition 停顿了事件时间一直卡在旧位置所有窗口状态都堆在内存里不触发。这时看执行计划看不出问题但要认识到窗口聚合状态清理完全依赖 watermark所以排查时一定要同步看当前 watermark 进度。3.3 去重专门去重器与配套状态Flink SQL 里的去重有好几种写法最常见的是ROW_NUMBER()去重和DISTINCT聚合。ROW_NUMBER去重典型写法SELECT * FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY ts ASC) AS rn FROM orders ) WHERE rn 1;物理计划里对应StreamExecRank如果优化器识别出去重语义会变成StreamExecDeduplicate。它的状态 key 是PARTITION BY的字段value 是“已经见过的记录标识”或排序序列。这里的状态大小取决于有多少条不重复的 key 在时间范围内没有被清理。COUNT(DISTINCT)这种聚合更隐蔽。它不会显示成独立算子而是会被优化器转成带状态的聚合内部需要一个 set 来保存所有已经出现过的去重值状态里每个 distinct 值都是集合里的一个元素。比如COUNT(DISTINCT user_id)状态量基本等于活跃用户数不是一个固定大小的计数。我在 1.1 里提到的那次事故就是这里吃了亏以为COUNT(DISTINCT)只是个计数实际上优化器在背后维护了一个巨大的去重集合。3.4 双流 JOIN左侧和右侧的缓冲区双流JOIN是另一个状态大户。看执行计划时只要看到StreamExecJoin基本就代表左右两侧都需要保留数据。因为流式 JOIN 是一个“缓存-配对”的过程一条数据进去之后不能马上决定输出得等另一边匹配的数据出现所以双方都得把历史数据放进状态。如果是带时间区间的 JOIN状态可以被区间边界约束住老数据会过期淘汰如果是不带时间条件的普通 JOIN那状态只增不减每一条历史数据都会长期驻留。最坑的是LEFT JOIN。业务方经常为了保留左表未匹配数据而写LEFT JOIN结果未匹配的那一侧数据会在状态里越堆越多而且没有天然清理机制。排查这种问题时执行计划里看到的还是同一个StreamExecJoin但状态增长模式完全不同。3.5 TopN 与排序排名状态TopN 常用的写法是SELECT * FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY category ORDER BY amount DESC) AS rn FROM orders ) WHERE rn 10;物理计划里对应StreamExecRank状态里保存的是当前每个分区的 TopN 记录。由于 N 是有限的状态量通常可控但如果分区键的基数特别大状态同样会涨。这里最忌讳的是写ORDER BY不带LIMIT对应StreamExecSort流式场景下相当于把整个数据流都攒在状态里做全局排序状态大小不受控。我见过有人把 OLAP 习惯带进 Flink SQL写了个全局ORDER BY的实时报表结果状态直接打满磁盘。3.6 维表 JOIN不是状态算子但别忽略缓存维表 JOIN 的物理算子是StreamExecLookupJoin它走的是查询外部存储MySQL、Redis、HBase的路线正常情况下不产生 Flink keyed state。真正影响内存的可能是维表缓存比如开启lookupcache 后本地缓存了最近查过的维表数据。这个缓存虽然不属于状态但也会吃掉堆内存而且因为不是严格的状态很多人会忽略它。总结成一页表就是SQL 类型典型物理算子状态 key状态 value清理语义分组聚合StreamExecGroupAggregate分组字段聚合中间结果靠 TTL无限增长窗口聚合StreamExecWindowAggregate分组字段窗口窗口聚合缓冲窗口触发后清理去重StreamExecRank / StreamExecDeduplicate去重字段已见记录集合靠 TTL双流 JOINStreamExecJoinJOIN key左右两侧历史数据无时间条件则无限增长TopNStreamExecRank分区字段TopN 记录集随新记录更新淘汰全局排序StreamExecSort无 key全量排序数据几乎不清理4. 用执行计划推定状态规模字段级估算与调优参数4.1 状态里一个 key 值值有多大看清了“哪个算子有状态”之后下一步要估“状态大概多大”。我一般会做一个粗略的字段级估算。拿StreamExecGroupAggregate(groupBy[user_id], select[user_id, COUNT(order_id), SUM(amount)])来说状态的 key 是user_idvalue 是聚合 buffer。COUNT内部是一个计数器SUM内部是一个数值累加器加起来可能不超过几十字节。但如果 value 里有LISTAGG、MAP、复杂的 POJO 结构序列化之后的字节数就需要按字段逐项估算。具体方法不复杂把状态 value 涉及的所有字段列出来按实际类型算序列化长度再乘以 key 数量得到总状态大小的数量级。这里必须强调这种估算只是数量级判断不用追求精确因为序列化框架有 header、对齐开销、压缩策略差异。目标是帮我们回答一个问题这个状态是“几百 MB”还是“几十 GB”完全不在一个量级。4.2 并行度怎么影响状态分布Flink 的 keyed state 是分区的数据按 key hash 到不同并行子任务每个子任务只保留自己负责的那部分 key。所以总状态量除以并行度大约就是单个 TaskManager 上的状态量。并行度越大单机压力越小但总状态量不会因为并行度变大而变小。很多人的误区是把并行度当成状态问题的手段。我一开始也是这么干的增加并行度只是把 6GB 状态从 3 个 TaskManager 分到了 6 个 TaskManager如果状态本身还在持续增长过几天又会打满新内存。并行度只能缓解短期压力真正要治的是状态本身的增长。4.3 状态后端的选择对规模的影响状态放在堆内存HashMapStateBackend里读写快但受限于堆大小GC 压力大状态放在 RocksDB 里序列化到磁盘单机可以放很大的状态但单次读写要过序列化和磁盘 IO。执行计划本身不会显示状态后端但同样的状态规模在这两种后端下表现完全不同几百 GB 的状态在堆内存里基本无解在 RocksDB 里还能撑。所以线上判断标准我一般是这样预估状态在单个 TaskManager 上低于 1GB用堆内存没问题超过 1GB 或者趋势是持续增长直接用 RocksDB再配本地磁盘。不要等到 Web UI 显示内存爆了再改。4.4 状态 TTL 是最后的保险绳不管是什么状态算子设置合理的table.exec.state.ttl都是最后一道保险。Flink 的状态清理是惰性的TTL 过期之后不会立刻删除而是等状态被访问或者后台清理线程来扫。我在生产里习惯给普通聚合状态设 1 小时左右给去重状态设 1 天以内具体看业务容忍度。真正常踩的坑是完全不设置 TTL或者设了一个“永远不过期”的时间。对于双流 JOIN、COUNT(DISTINCT)这种天然无限增长的状态没有 TTL 等于把作业变成一个慢慢被填满的内存炸弹。执行计划不会直接显示 TTL 配置所以要在 SQL 配置或者 CLI 里主动确认。5. 我查执行计划时积累的几个实操习惯5.1 上线前把 EXPLAIN 存档而不是等出问题再看现在我把“上线前存档 EXPLAIN”当成和测试用例一样的必做动作。每次提交新的 SQL我会在 SQL Client 跑一次EXPLAIN然后连同建表 DDL、拓扑结构一起截图存到当前任务文档里。一旦线上出问题拿出来对照能节省很多时间。尤其是团队协作时别人接手这个作业第一件事也是看执行计划确认状态在哪而不是猜业务逻辑。5.2 重点扫三类异常信号拿到一份执行计划我不逐行细读而是扫三个信号hash分布后面跟了带状态算子说明这里一定有 keyed state要去确认 key 基数Join节点左右两边有没有时间过滤条件没有的话基本等于无限状态Rank或者Deduplicate前面如果有全局hash注意确认分区键是不是高基数。这三个信号扫完状态大局基本就清楚了。5.3 别忽略优化器做的“隐形改写”有时候你写的 SQL 和自己以为的完全不是一回事。比如COUNT(DISTINCT)会被改写成带 set 语义的聚合LEFT JOIN在某些情况下会保留更多状态ROW_NUMBER去重和 TopN 在物理计划里可能是同一个算子但状态语义完全不同。优化器改写不是坏事它通常让 SQL 跑得更高效。但如果你只看业务 SQL 不看执行计划就会对真实状态一无所知。我见过有人写COUNT(DISTINCT)满脑子以为自己算的是一个数字实际上状态里已经攒了几百万个去重值这就是不看执行计划的代价。5.4 用 Web UI 去验证而不是猜最后一个小建议EXPLAIN帮你做定性判断Web UI 帮你做定量验证。作业跑起来之后去 Web UI 的每个算子详情页看 “State” 相关指标能看到每个算子的状态大小、状态读写次数。把这些数据和执行计划里的算子一一对应起来就能确认到底是哪一个算子出了问题而不是盯着整个作业的 checkpoint 大小发愁。我现在处理状态类问题的固定套路就是SQL 导出 - 看EXPLAIN找状态算子 - 估算状态量级 - 检查 TTL - 用 Web UI 指标定位具体 TaskManager。这套链路跑下来绝大多数状态膨胀问题都能在两小时内定位到算子层面。状态本身不是洪水猛兽每个实时计算作业多多少少都要用状态。真正让人头疼的是“不知道状态在哪、有多大、什么时候清理”。EXPLAIN PLAN就是把这个“不知道”变成“知道”的最快路径。所以如果你下次遇到状态问题别急着加内存先花十分钟把执行计划打开你会发现答案早就写在里面了。