You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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自动生成的类

我为每类数据单独编写了流处理作业,逻辑为通过时间窗口缓存元素并聚合为列表,每类作业按统一流程测试:

  1. 启动作业正常运行
  2. 触发savepoint后停止作业
  3. 为对应的数据类新增一个字段
  4. 基于之前生成的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正常恢复:

  1. 先把POJO类改到符合Flink规则
    给io.peleg.pojo.User补上@NoArgsConstructor注解,保证类是public修饰,所有字段要么public,要么有符合JavaBean规范的getter/setter;后续新增字段必须加@Nullable注解,反序列化旧状态时Flink会自动给缺失字段赋值null,不会报错。
  2. 所有算子显式声明返回类型,彻底堵死类型擦除问题
    不要依赖Flink自动类型推导,实现WindowFunction/AggregateFunction时重写getProducedType()方法返回对应TypeInformation;如果用returns()方法声明类型,不要写returns(List.class)这种擦除泛型的写法,用TypeHint指定完整泛型:
    .returns(new TypeHint<List<User>>() {})
    
  3. Avro类型显式绑定Avro序列化器
    对Avro生成的User类,直接用AvroTypeInfo生成TypeInformation绑定到对应算子,不要让Flink自动推导。
  4. 本地测试阶段开配置提前暴露序列化问题
    本地调试时加上这行配置:
    env.getConfig().disableGenericTypes();
    
    开启后只要Flink要回退到Kryo就直接抛异常,不用等跑savepoint恢复时才发现序列化器用错,排查效率高很多。
  5. 固定Flink镜像版本和本地依赖对齐
    把compose里的镜像从flink:latest改成flink:1.15.0,保证集群和作业依赖版本完全一致,避免版本差异导致的兼容问题。

注意:Kryo序列化的状态天生不支持Schema演化,只要状态存储走了Kryo,修改类结构后就不可能从旧savepoint恢复,必须保证所有状态存储的类型都命中Flink自带的支持演化的序列化器(POJO、Avro、Tuple、Value类型等)。

内容的提问来源于stack exchange,提问作者Peleg Tsadok

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 03:36:25