ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink StateMigrationException排查与状态迁移实战指南

Flink StateMigrationException排查与状态迁移实战指南 1. 一场升级引发的线上事故先看报错现场凌晨两点值班手机把我从梦里拽出来。告警说的是线上一个Flink SQL作业连续重启失败作业已经进入FAILED状态。登录平台一看日志罪魁祸首是这一行Caused by: org.apache.flink.util.StateMigrationException: The new state serializer must not be incompatible说句实话看到StateMigrationException的那一刻我反而松了一口气。这算是Flink状态恢复里比较“有名”的异常之一虽然名字吓人但大多数情况下并不是什么底层bug而是作业自身的状态结构发生了不兼容变化。真正麻烦的是后面那句“The new state serializer must not be incompatible”字面意思是“新的状态序列化器不能是不兼容的”读起来有点绕翻译成人话就是Flink在从旧状态恢复时发现算子要用的新序列化器跟之前写入状态时用的序列化器对不上校验直接失败。这个作业本身很简单就是从Kafka读数据做一次分组聚合再写回下游。改动听起来也不大——业务方说“就加了一个维度字段”我按照需求改了SQL里的GROUP BY然后触发了一次从最近savepoint恢复的升级。结果就是上面那个异常作业卡在恢复阶段起不来。这个报错最折磨人的地方在于它不是每次必现而是跟状态里实际保存的内容有关系。状态里有旧数据新的序列化器跟旧的描述对不上Flink宁可失败也不冒险因为一旦状态被错误反序列化轻则数据算错重则状态直接损坏。所以从设计上看Flink这个校验是保护机制不是误报。这次事故从发生到彻底解决前后折腾了大半天。中间踩了不少坑也顺手把Flink状态迁移相关的机制重新捋了一遍。这篇文章就把整个排查思路、背后原理和实际方案完整记录下来给遇到同样问题的朋友一个参考。2. 扒开StateMigrationException的外衣核心机制与出错原因2.1 Flink状态恢复时到底在做什么要理解这个异常得先搞清楚Flink从savepoint或checkpoint恢复时的基本流程。每个有状态算子在初始化时都会创建StateDescriptor描述符里带着状态名、状态类型ValueState、ListState、MapState这种和对应的StateSerializer。当作业启动并尝试从外部存储恢复状态时Flink会将当前算子要创建的描述符与状态存储里记录的历史描述符做一次比对。比对的核心逻辑在AbstractStateBackend的代码里简单说就是三步检查状态名是否存在、检查状态类型是否一致、检查序列化器是否兼容。前两步相对好理解状态名对不上说明拓扑结构变了可能漏了某个算子或者改了算子ID状态类型变了比如从ValueState变成了MapState那肯定恢复不了。真正容易出问题的就是第三步——序列化器的兼容性校验。兼容性校验的源码逻辑大致是这样的if (previousSerializer ! null !newSerializer.getCompatibleSerializer().getClass().equals(previousSerializer.getClass())) { throw new StateMigrationException(The new state serializer must not be incompatible); }这里有个关键点Flink并不是要求新旧序列化器“完全相等”而是要求新序列化器的类型必须跟旧序列化器的类型保持一致。同时它还允许通过StateMigrationContext做数据迁移。也就是说如果新序列化器的类型不同即使字段含义完全一样也会抛出StateMigrationException。2.2 为什么SQL作业特别容易触发这个问题Flink SQL作业有一个特性状态结构不是我们手动定义的而是由Flink SQL框架根据SQL语义自动推算出来的。这套推算对用户是黑盒用户看的是SQL逻辑但Flink底层会生成一串算子链每个算子里面有它自己定义的状态。SQL层面的很多“小改动”映射到底层状态结构上就是“大变化”。我遇到的情况就是一个典型例子。旧SQL大致长这样CREATE TABLE source_table ( user_id BIGINT, category_id BIGINT, action_time TIMESTAMP(3), action_type STRING, WATERMARK FOR action_time AS action_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic source_topic, properties.bootstrap.servers kafka:9092, properties.group.id source_group, format json, scan.startup.mode latest-offset ); CREATE TABLE sink_table ( user_id BIGINT, cnt BIGINT ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/flink_test, table-name sink_table, username root, password 123456 ); INSERT INTO sink_table SELECT user_id, COUNT(*) FROM source_table GROUP BY user_id;新SQL改动是增加一个统计维度GROUP BY后面多了一个字段INSERT INTO sink_table SELECT user_id, category_id, COUNT(*) FROM source_table GROUP BY user_id, category_id;就这一个小改动底层状态结构发生了三重变化第一聚合算子的Key类型变了。原来GROUP BY只有一个BIGINT字段Key是BinaryRowData包装的单个字段现在有两个字段Key变成了复合结构对应的TypeSerializer不再是同一个实现类。第二聚合结果的状态描述符变了。SQL聚合算子底层用的是MapStatekey是命名空间加逻辑key的组合value是累加器。字段数量变化会导致累加器Accumulator的内部结构变化进而导致状态value的序列化器type不同。第三算子的UID可能都变了。Flink SQL生成的算子在底层是有编号的如果拓扑diff比较大某些算子可能被替换或重新生成状态名对不上就直接报错了。2.3 除SQL改动之外的高频触发场景除了我这次遇到的“改GROUP BY字段”之外还有几个常见场景也容易触发StateMigrationException我把它们一并列出来Flink版本升级。这是最典型的场景比如从Flink 1.13升到1.17中间Flink SQL框架对聚合状态的内部表示做过多次调整。特别是涉及窗口聚合、会话窗口这类复杂算子时状态序列化器不兼容的情况非常普遍。网上大量StateMigrationException的帖子都出现在大版本升级之后不是说代码写错了纯粹是框架内部序列化结构变了。维表JOIN缓存状态变化。Flink SQL维表JOIN的LookupJoin算子会为每个流式侧数据创建一个临时状态用于记录上次查询结果的过期时间配合TTL使用。如果维表的主键类型变更、JOIN条件变化或者缓存配置方式变了同样会导致状态描述符比对失败。UDF返回类型变更。自定义UDF的返回类型从一种变成了另一种如果这个UDF被用在聚合函数里累加器里保持的中间结果类型就会跟着变。看起来只是改了一个函数的modifier实际上聚合状态的序列化器已经完全不同。从Raw State切到Managed State。这属于比较进阶的问题。如果历史作业里用了自定义的KeyedStateBackend直接写入原始状态后来重构代码改用托管状态那么之前写的原生状态因为没有 serializer元数据Flink无法正确做兼容性判断启动时大概率直接抛异常。State TTL设置的变更。严格说TTL变更一般不直接导致serializer不兼容但为了让TTL生效Flink会在状态里包一层TtlStateFactory里面会引入额外的序列化器包装逻辑。某些版本下新旧TTL配置切换时可能与序列化器结合方式变化有关踩过坑的朋友应该懂我在说什么。了解这些触发场景很重要因为不同原因对应的处理策略完全不同。如果直接照着一个方案生搬硬套后面大概率还会再出事。3. 从根源到对策我的排查思路和解决路径3.1 第一步先把“元凶”找出来遇到这个异常不要急着改代码或清状态。先搞清楚三件事哪个算子报错、状态名是什么、新旧作业的拓扑差异在哪里。这里有个小技巧看日志不能只看最后一行异常要往上翻。在StateMigrationException之前通常会有一段类似这样的日志2025-01-15 02:15:32,678 INFO org.apache.flink.runtime.state.heap.HeapKeyedStateBackend - Trying to restore keyed state from savepoint 2025-01-15 02:15:32,679 INFO org.apache.flink.runtime.state.heap.HeapKeyedStateBackend - Trying to restore state of operator 2f2e3a1d9e0b4c5f8a6b7c8d9e0f1a2b注意在“Trying to restore keyed state from savepoint”之后的那几行日志里面带着operator的十六进制哈希ID。我们要找的就是这个算子ID它可以帮我们在没出问题之前先去作业拓扑里对应算子。如果平台上能看到UDF算子名或算子链信息那就更好办了。找到报错算子后用两个作业的JSON plan做对比。具体操作是用flink sql -f old.sql时的Web UI页面导出旧作业的JobGraph JSON修改后的新SQL本地也编译一份再用diff工具对比。状态描述符的信息就藏在JobGraph里。如果环境里不方便导出JSON还有一个更“笨”但有效的办法逐个回退SQL修改点。新版SQL是一次改多了就先把改动拆开每次只动一处试恢复。这个方法在正在排障时很实用因为直接从“变的点”去定位比对着源码分析快得多。3.2 从日志锁定问题算子和状态信息我当时的排查过程就是顺着日志找。报错日志里除了顶部的异常摘要还会打印当前算子的状态描述符信息例如Caused by: org.apache.flink.util.StateMigrationException: The new state serializer must not be incompatible at org.apache.flink.runtime.state.KeyedStateBackend.verifyStateMigration(...) at org.apache.flink.runtime.state.heap.HeapKeyedStateBackend.createInternalState(...) at org.apache.flink.streaming.api.operators.AbstractKeyedStateBackend.createState(...)接着下面有类似这样的描述Unable to restore state of operator 2f2e3a1d9e0b4c5f8a6b7c8d9e0f1a2b because the new serializer is not compatible with the previous state.这个operator ID就是定位的关键。当时我操作平台的界面里能先点开作业拓扑看到每个算子但节点名字都是“GroupAggregate”、_group_agg这类框架生成的名称看不出哪个对应具体状态。后来我用Flink SQL自己本地提交了一份新版作业草稿通过state.backend配置让作业停在“RUNNING”前的初始化阶段再把JobGraph export出来跟平台上的旧作业JobGraph做diff。果然聚合算子的operator ID发生了变化。原因就是新增GROUP BY字段后SQL优化器生成了一个新的group aggregate算子ID自然对不上。状态名直接从“旧名字”找不到Flink就只能抛出异常了。所以当聚合字段变化时本质上不是serializer不兼容的问题而是状态整体找不到的问题——这里Serialization不兼容的时候是因为状态名能对上但新的serializer类和旧的serializer类对不上了。这两类信息日志里都会明明白白地告诉你。看的时候别先看异常栈要先看异常信息里有没有出现状态名或字段类型的关键字。3.3 第二步判断当前场景能不能“无损”解决定位到原因后思路立刻清晰了很多。现在的问题不是“为什么失败”而是“想不想保留旧状态”。这两个场景要分开处理场景A旧状态不能丢。这种情况最常见于生产环境有实时指标、有累计结果、有中间状态需要衔接的场景。比如统计昨天的UV不能因为作业升级就把计数器清零重新算。又比如用Flink做了几个月的窗口聚合中间状态里有大量未结束的会话丢了就等于丢掉业务数据。这种时候不能清状态重启必须想办法迁移或者让状态结构保持一致。场景B旧状态可以丢。如果作业是纯粹的清洗转发或者统计任务对“重启后从零开始”这件事可以接受那就简单多了。把状态清了再恢复是最省事也最不容易出错的办法。很多测试环境、辅助分析作业都是这个情况。我的这个作业属于场景A——线上统计不允许丢状态。所以重点放在“如何让新serializer兼容旧状态”上。3.4 方案A让SQL改动不改变状态结构保状态首选既然新增GROUP BY字段会导致聚合状态结构完全变化那么从逻辑上绕过去不要改动GROUP BY层级把新增字段放到下游别的计算环节中去处理。比如原来的统计粒度是user_id新版需要user_id category_id两个维度。那完全可以保留原有SQL只做user_id级别聚合增加一个新的聚合算子处理更细粒度或者通过多级关联去生成细粒度结果。但问题是如果业务方明确需要“细粒度结果直接写入sink”这种绕路方式会让SQL变复杂很多而且原作业的语义就变了。还有一种更常见的写法变更是“改SELECT字段顺序”。比如原来输出是SELECT user_id, cnt现在改成SELECT cnt, user_id这个改动本身不影响GROUP BY的key但会影响Sink端的字段映射。如果sink连接器的schema是强类型匹配字段顺序变了反而可能引入新的序列化错误但一般不会触发StateMigrationException。因为它动的不是状态逻辑。真正能保持状态结构不变的是那种看起来改了SQL、实际上底层拓扑不变或者状态名不变的情况修改过滤掉的数据范围比如WHERE条件从action_type click改成action_type IN (click, view)。这通常对聚合状态结构无影响。修改输出字段的投影但保留聚合的输入字段不变。增加多余的UDF只要它的输入输出类型不变通常也不会改动状态结构。所以方案A的核心就是一个词还原。回到报错前的状态把SQL改动拆细找到真正导致状态变化的那个点然后用一种“不改变聚合算子结构”的方式去满足新需求。这次我的处理就是这一条路。我没有直接改原SQL的GROUP BY而是保留原聚合SQL不变在外面套一层子查询做gmv重命名和下游维度补充从而在不触碰聚合状态的情况下实现了新增维度。3.5 方案B接受丢状态无状态启动最快恢复如果场景B成立那方案就非常直接了。操作上就是做一次无状态恢复。Flink从savepoint恢复时可以通过--allowNonRestoredState忽略不上的状态。但要注意——这里有一个常见的坑allowNonRestoredState只是允许跳过找不到的状态真正可以做到完全不恢复任何状态是用一个干净的启动而不是从旧savepoint恢复。最干净的做法是根本无法从失败的作业的savepoint直接恢复而是直接不指定savepoint启动新作业。如果作业本身没有强状态依赖直接不带savepoint重启即可。如果你的平台允许“从空状态启动”那是最好的。把旧作业stop掉新作业以无状态方式启动。Kafka源从latest-offset或指定时间点读取任务就能跑起来。需要注意如果Kafka里积压了大量未消费的数据无状态启动可能会让部分中间统计失效比如重启窗口计算时重启前的数据窗口会被丢弃在窗口聚合统计的场景里这个要提前跟业务确认。如果一定要用旧savepoint恢复同时想让Flink忽略掉那一部分对不上的状态你可以加这个启动参数--allowNonRestoredState在SqlClient或运维提交界面上一般对应一个allowNonRestoredState开关。找到那个勾选框打开也能达到同样效果。不过我强烈建议凡是打开这个开关前都要先搞清楚被忽略的到底是什么状态。因为如果漏掉的正是核心聚合状态那数据从恢复那一刻起就算不出来了运行期间不会报错直到你发现指标差了一大截。这就比启动失败更可怕了。3.6 方案C用State Processor API做状态迁移高级玩法如果方案A和B都不合适——既不能丢状态又不能通过SQL结构保持兼容——那就只剩最后一条路主动做状态迁移。Flink提供了一个专门用来“离线读写外部状态”的API叫State Processor API。它能让你在不启动作业的情况下把savepoint里的状态读出来做各种转换再写回一个新的savepoint。核心代码逻辑大概是SavepointReader savepointReader SavepointReader.read(env, hdfs:///flink/savepoints/savepoint-xxx); DataStreamKeyedStateReaderFunctionString, AccumulatorState stateStream savepointReader.readKeyedState(_agg_operator_uid, new KeyedStateReaderFunctionString, AccumulatorState() { Override public void readKey(String key, Context ctx, CollectorAccumulatorState out) throws Exception { // 读取旧状态的value值 } });然后定义新的状态写入SavepointWriter.newSavepoint(env, flink-1.17).withKeyedState( new KeyedStateWriterFunctionString, AccumulatorState() { Override public void writeKey(String key, Context ctx, AccumulatorState value) throws Exception { // 写入新结构的MapState } } ).write(hdfs:///flink/savepoints/savepoint-migrated);这里有个前置条件旧作业必须是用托管状态生成的不是Raw State。如果是Raw StateState Processor API是读不到状态的因为根本没有序列化器元数据在里面。这个方案最大的价值在于它允许你对状态做结构转换比如把旧聚合状态里的累加器数据“拆”到新结构里。缺点同样明显代码量不算小需要对Flink状态API本身有一定掌握调试成本不低。实际用不用这个方案取决于“状态迁移”这件事值不值得付出这些成本。很多团队其实在绝大多数情况下根本用不上State Processor API因为SQL作业的状态结构是框架内部生成的你读出来的Accumulator对象本身可能依赖于Flink内部类框架版本升级后连对象反序列化都有问题。所以除非团队里有比较熟悉底层的同学否则我一般不推荐一上来就搞这个。3.7 方案D治理思路——规范UID和兼容性治理说到底StateMigrationException最好的结局是“根本不出现”。想要不出现就必须在作业开发和发布流程上做一些规范。**第一个规范给关键算子手动指定UID。**看似是小事但作用很大。手动指定UID后即使SQL逻辑微调只要UID没变Flink就仍会把旧状态映射到对应的算子上配合序列化器兼容性校验成功率会大大提升。**第二个规范对状态结构改动做影响评估。**在改动任何一个涉及GROUP BY、JOIN、窗口或状态UDF的SQL之前先问一句“这次改动会影响聚合算子底层的状态结构吗”。如果会就主动设计状态迁移方案而不是等发布后报错再救火。**第三个规范保存两个版本的savepoint。**升级前先手动触发一个savepoint升级失败后再尝试从旧的savepoint回滚这个操作能帮你争取不少时间。很多平台会自动清理旧savepoint但这不等于说没问题。像这类状态兼容问题一个重要恢复策略是你手头还有一份能用的、与旧作业完全匹配的savepoint。我在这次事故里其实就吃了没有提前做“手动UID”的亏。因为旧作业的SQL完全是由平台自动生成的算子UID结构调整后一堆算子ID都变了连排查都比正常情况更费劲。4. 实践中的坑与心得处理状态兼容问题的避坑指南4.1 SQL变更兼容性清单根据我这次踩坑的经验下面这份“SQL变更→状态影响”对照表建议大家收藏起来。每次改动完SQL先对照检查一遍。变更类型是否影响状态结构风险等级说明修改SELECT输出字段不改GROUP BY不影响状态结构低聚合Key不变序列化器不变修改WHERE条件/过滤逻辑通常不影响状态结构低但窗口数据范围会变GROUP BY字段增加/减少影响Key类型高聚合算子状态结构整体变化JOIN条件变化影响Key类型高LookupJoin缓存状态受影响窗口类型/步长变化影响状态结构高窗口聚合状态内部结构变化UDF返回类型变化可能影响累加器中需专门验证增加/删除聚合函数影响累加器字段中先本地验证再操作Flink版本升级影响所有序列化器类高排在所有SQL变更风险之上我在实际项目中发现最容易出问题的是“三个高”叠加——高复杂度SQL多层子查询、高频率迭代一周改三次、高位环境依赖平台自动生成UID。三个条件同时满足时StateMigrationException几乎必然会发生。所以在设计SQL上线流程时这些风险点要尽可能前置。4.2 排查状态问题的实操命令和平台操作排查过程中的几个实用命令也一并分享出来# 从savepoint恢复作业带允许忽略状态开关 flink run -s hdfs:///flink/savepoints/savepoint-xxx \ -d \ -c com.example.MyJob \ --allowNonRestoredState \ ./my-job.jar如果使用的是Flink SQL作业很多平台提交入口在界面上有“恢复路径”和“允许忽略状态”两个选项直接勾上就行。但再次强调这个开关是逃生通道不是常规手段。还有一个从日志里定位状态的方法。如果你的作业已经暴露了Web UI打开TaskManager的日志页面搜索关键字的上下文状态恢复失败时会打印出所有待恢复的状态描述符信息。拿到状态名后再用StateProcessorAPI或FLIP-81的兼容性检查工具去验证。4.3 TTL和状态迁移的联动坑这个坑我差点踩翻。状态TTLTime-To-Live的配置改动也会影响恢复时的行为。有些状态下旧作业没有配置TTL新作业配置了TTLFlink会尝试给旧状态套上一层TTL包装。如果包装逻辑和序列化器组合出现版本兼容问题就可能出现“连状态都找到了但serializer还是对不上”的情况。我的建议是**状态TTL的变更尽量跟状态结构的变更分开进行先单独升级TTL并验证成功再动状态结构。**虽然这会多占一次发布窗口但每一步都可控。否则两个变量叠加在一起出错时你很难定位到底是哪一刻引入的问题。4.4 一个被低估的“土办法”双跑校验可能有些朋友觉得前面讲的方案都好麻烦有没有更省事的有就是做双跑验证。做法很简单新版本作业在测试环境运行与生产环境的旧版本并行一段足够长的时间比较两者产出的数据。如果指标对得上说明状态结构的事实是OK的。测试通过后择机切换从最新checkpoint恢复注意不是旧savepoint。这个做法在时间成本上可能更高但确实能规避绝大多数状态兼容性问题尤其是SQL改动加版本升级一起来的时候。而且它还能顺带验证“改动后计算逻辑正确性”一箭双雕。4.5 一点长期经验状态治理要前置处理这次问题后的最大体会是状态兼容问题不是一个“技术单点问题”而是作业开发流程和治理机制一起要面对的问题。如果团队里有Flink SQL作业的迭代需求应该在设计阶段就把状态兼容考虑进去。比如尽量手动给state命名和指定UID不要完全依赖系统自动生成再比如建立StateMigrationTest回归用例发布前自动跑一遍状态兼容性检查。这块投入的成本不高但能省掉后面每个凌晨的救火时间。另外云上环境如果支持可以用作业快照Snapshots配合状态过期策略做状态版本管理给每次发布自动打一个“可回滚”的状态备份。一旦升级失败一键回滚。这套机制值得研究因为它解决的正是“升级失败无法回滚”这个最痛的点。5. 结尾说点实际操作后的真心话这段经历之后我对StateMigrationException有了完全不同的理解。最初遇到这个异常时我第一反应是Flink的bug想着怎么绕过、怎么清状态。但捋清楚之后发现这更像是Flink在“温柔地”提醒我你改了状态结构但是状态里还有旧数据直接起跑我办不到。现在我的处理优先级非常固定先看能否还原SQL、保持状态结构不变再看能否接受清状态重启如果都不行就提前规划State Processor API或者设计双跑验证。这四步走下来这个报错的处理效率高了很多至少不会再凌晨爬起来一脸懵。最后再分享一个小细节如果你的作业有状态迁移需求最好在开发环境提前把新旧两个版本跑一遍对比savepoint里各个状态名和序列化器描述。别等到生产环境出了异常才开始学习状态管理的原理。状态这东西就像是Flink作业的“记忆”。你可以在它遗忘时重建也可以在它改变时迁移——但前提是你得知道它在哪、长什么样、该怎么动。搞明白这一层StateMigrationException就不再是玄学。
返回列表