Flink对POJO/Avro类默认使用Kryo序列化致状态Schema演化失败
问题背景
我在Flink 1.15.0、Java 11环境下开展State Schema Evolution(Flink状态模式演化)POC验证,针对三类序列化场景分别创建了对应的数据类:
io.peleg.kryo.User:包含java.time.Instant类型字段,已知该类不满足Flink POJO序列化要求io.peleg.pojo.User:仅包含Integer、Long、String这类基础包装类型字段,getter、setter、构造方法均由Lombok生成io.peleg.avro.User:通过Avro Maven插件根据Avro Schema自动生成的类
我为每类数据单独编写了流处理作业,逻辑为通过时间窗口缓存元素并聚合为列表,每类作业按统一流程测试:
- 启动作业正常运行
- 触发savepoint后停止作业
- 为对应的数据类新增一个字段
- 基于之前生成的savepoint提交作业尝试恢复
异常现象
三类作业从savepoint恢复时全部抛出相同异常,核心报错栈显示状态反序列化时走了Kryo序列化器,抛出索引越界:
java.lang.Exception: Exception while creating StreamOperatorStateContext. at org.apache.flink.streaming.api.operators.StreamTaskStateInitializerImpl.streamOperatorStateContext(StreamTaskStateInitializerImpl.java:255) at org.apache.flink.streaming.api.operators.AbstractStreamOperator.initializeState(AbstractStreamOperator.java:268) at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.initializeStateAndOpenOperators(RegularOperatorChain.java:106) at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreGates(StreamTask.java:700) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.call(StreamTaskActionExecutor.java:55) at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreInternal(StreamTask.java:676) at org.apache.flink.streaming.runtime.tasks.StreamTask.restore(StreamTask.java:643) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:948) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:917) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:741) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:563) at java.base/java.lang.Thread.run(Unknown Source) Caused by: org.apache.flink.util.FlinkException: Could not restore keyed state backend for WindowOperator_3983d6bb2f0a45b638461bc99138f22f_(2/2) from any of the 1 provided restore options. at org.apache.flink.streaming.api.operators.BackendRestorerProcedure.createAndRestore(BackendRestorerProcedure.java:160) at org.apache.flink.streaming.api.operators.StreamTaskStateInitializerImpl.keyedStatedBackend(StreamTaskStateInitializerImpl.java:346) at org.apache.flink.streaming.api.operators.StreamTaskStateInitializerImpl.streamOperatorStateContext(StreamTaskStateInitializerImpl.java:164) ... 11 more Caused by: org.apache.flink.runtime.state.BackendBuildingException: Failed when trying to restore heap backend at org.apache.flink.runtime.state.heap.HeapKeyedStateBackendBuilder.restoreState(HeapKeyedStateBackendBuilder.java:172) at org.apache.flink.runtime.state.heap.HeapKeyedStateBackendBuilder.build(HeapKeyedStateBackendBuilder.java:106) at org.apache.flink.runtime.state.hashmap.HashMapStateBackend.createKeyedStateBackend(HashMapStateBackend.java:143) at org.apache.flink.runtime.state.hashmap.HashMapStateBackend.createKeyedStateBackend(HashMapStateBackend.java:74) at org.apache.flink.runtime.state.StateBackend.createKeyedStateBackend(StateBackend.java:140) at org.apache.flink.streaming.api.operators.StreamTaskStateInitializerImpl.lambda$keyedStatedBackend$1(StreamTaskStateInitializerImpl.java:329) at org.apache.flink.streaming.api.operators.BackendRestorerProcedure.attemptCreateAndRestore(BackendRestorerProcedure.java:168) at org.apache.flink.streaming.api.operators.BackendRestorerProcedure.createAndRestore(BackendRestorerProcedure.java:135) ... 13 more Caused by: com.esotericsoftware.kryo.KryoException: java.lang.IndexOutOfBoundsException: Index 83 out of bounds for length 3 Serialization trace: favoriteColor (io.peleg.avro.User) at com.esotericsoftware.kryo.serializers.ObjectField.read(ObjectField.java:125) at com.esotericsoftware.kryo.serializers.FieldSerializer.read(FieldSerializer.java:528) at com.esotericsoftware.kryo.Kryo.readClassAndObject(Kryo.java:761) at com.esotericsoftware.kryo.serializers.CollectionSerializer.read(CollectionSerializer.java:116) at com.esotericsoftware.kryo.serializers.CollectionSerializer.read(CollectionSerializer.java:22) at com.esotericsoftware.kryo.Kryo.readClassAndObject(Kryo.java:761) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.deserialize(KryoSerializer.java:402) at org.apache.flink.runtime.state.heap.HeapSavepointRestoreOperation.readKVStateData(HeapSavepointRestoreOperation.java:219) at org.apache.flink.runtime.state.heap.HeapSavepointRestoreOperation.readKeyGroupStateData(HeapSavepointRestoreOperation.java:149) at org.apache.flink.runtime.state.heap.HeapSavepointRestoreOperation.restore(HeapSavepointRestoreOperation.java:125) at org.apache.flink.runtime.state.heap.HeapSavepointRestoreOperation.restore(HeapSavepointRestoreOperation.java:57) at org.apache.flink.runtime.state.heap.HeapKeyedStateBackendBuilder.restoreState(HeapKeyedStateBackendBuilder.java:169) ... 20 more Caused by: java.lang.IndexOutOfBoundsException: Index 83 out of bounds for length 3 at java.base/jdk.internal.util.Preconditions.outOfBounds(Unknown Source) at java.base/jdk.internal.util.Preconditions.outOfBoundsCheckIndex(Unknown Source) at java.base/jdk.internal.util.Preconditions.checkIndex(Unknown Source) at java.base/java.util.Objects.checkIndex(Unknown Source) at java.base/java.util.ArrayList.get(Unknown Source) at com.esotericsoftware.kryo.util.MapReferenceResolver.getReadObject(MapReferenceResolver.java:42) at com.esotericsoftware.kryo.Kryo.readReferenceOrNull(Kryo.java:805) at com.esotericsoftware.kryo.Kryo.readObjectOrNull(Kryo.java:728) at com.esotericsoftware.kryo.serializers.ObjectField.read(ObjectField.java:113) ... 31 more
原本预期只有io.peleg.kryo.User对应的作业会恢复失败——毕竟Kryo本身不支持Schema演化,但实际测试发现三类数据全部没有命中POJO、Avro序列化器,全部回退到了Kryo序列化。
测试用Flink集群通过docker compose部署,配置如下:
version: "2.2" services: jobmanager: image: flink:latest ports: - "8081:8081" command: jobmanager environment: - | FLINK_PROPERTIES= jobmanager.rpc.address: jobmanager taskmanager: image: flink:latest depends_on: - jobmanager command: taskmanager scale: 1 environment: - | FLINK_PROPERTIES= jobmanager.rpc.address: jobmanager taskmanager.numberOfTaskSlots: 2
测试目标:实现POJO类增删字段后,作业可基于旧版本生成的savepoint正常启动恢复。
根因分析
所有类全部回退到Kryo序列化,本质是Flink类型推导失败,常见触发点按概率排序:
- 窗口聚合的泛型返回值被类型擦除
作业逻辑是窗口内攒批返回List<User>,Java泛型在运行时会擦除类型信息,如果没有给算子显式声明返回值类型,Flink根本不知道列表里存的是什么对象,直接回退到通用Kryo序列化整个列表——这也是连Avro生成的User类都走Kryo的核心原因,不是User类本身的序列化器没选对,是外层List直接被Kryo兜底了。 - Lombok生成的POJO不满足Flink POJO判定规则
Flink识别POJO的硬要求是类必须有public无参构造,很多人用Lombok时只加@Data、@AllArgsConstructor,漏掉@NoArgsConstructor,这种情况Flink不会把类识别为支持Schema演化的POJO类型,直接回退Kryo。除此之外类是final修饰、字段没有标准getter/setter也会判定失败。 - Avro类型未显式绑定专用序列化器
Flink不会自动识别Avro生成类并匹配Avro序列化器,必须显式指定AvroTypeInformation,否则类型推导阶段只会把它当成普通Java类,最终还是走Kryo。 - Docker镜像版本和本地依赖不一致
Compose配置里用的是flink:latest镜像,本地作业依赖是Flink 1.15.0,一旦latest拉到更高版本,版本差异会导致序列化器注册、状态兼容逻辑异常,直接触发序列化回退。
修复方案
按以下顺序调整,即可实现POJO增删字段后从savepoint正常恢复:
- 先把POJO类改到符合Flink规则
给io.peleg.pojo.User补上@NoArgsConstructor注解,保证类是public修饰,所有字段要么public,要么有符合JavaBean规范的getter/setter;后续新增字段必须加@Nullable注解,反序列化旧状态时Flink会自动给缺失字段赋值null,不会报错。 - 所有算子显式声明返回类型,彻底堵死类型擦除问题
不要依赖Flink自动类型推导,实现WindowFunction/AggregateFunction时重写getProducedType()方法返回对应TypeInformation;如果用returns()方法声明类型,不要写returns(List.class)这种擦除泛型的写法,用TypeHint指定完整泛型:.returns(new TypeHint<List<User>>() {}) - Avro类型显式绑定Avro序列化器
对Avro生成的User类,直接用AvroTypeInfo生成TypeInformation绑定到对应算子,不要让Flink自动推导。 - 本地测试阶段开配置提前暴露序列化问题
本地调试时加上这行配置:
开启后只要Flink要回退到Kryo就直接抛异常,不用等跑savepoint恢复时才发现序列化器用错,排查效率高很多。env.getConfig().disableGenericTypes(); - 固定Flink镜像版本和本地依赖对齐
把compose里的镜像从flink:latest改成flink:1.15.0,保证集群和作业依赖版本完全一致,避免版本差异导致的兼容问题。
注意:Kryo序列化的状态天生不支持Schema演化,只要状态存储走了Kryo,修改类结构后就不可能从旧savepoint恢复,必须保证所有状态存储的类型都命中Flink自带的支持演化的序列化器(POJO、Avro、Tuple、Value类型等)。
内容的提问来源于stack exchange,提问作者Peleg Tsadok
相关产品推荐
相关产品推荐

