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

Flink 1.15在AWS KDA中恢复Savepoint失败的问题求助

我在AWS KDA(托管Apache Flink)上运行Flink 1.15应用,更新时执行Snapshot(savepoint)操作后恢复失败,推测原因是迁移了ValueState中封装类的路径——将old.class.path.MyClass迁移到了new.class.path.MyClass。

错误日志

threadName (7eedaf4db54f14ae5b70ce55bfe409f0) switched from INITIALIZING to FAILED with failure cause: 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:953)

    at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:922)

    at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:746)

    at org.apache.flink.runtime.taskmanager.Task.run(Task.java:568)

    at java.base/java.lang.Thread.run(Thread.java:829)

Caused by: org.apache.flink.util.FlinkException: Could not restore keyed state backend for KeyedProcessOperator_5856e2d12487eeccb905e376fa806891_(10/32) 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: Caught unexpected exception.

    at org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackendBuilder.build(RocksDBKeyedStateBackendBuilder.java:395)

    at org.apache.flink.contrib.streaming.state.EmbeddedRocksDBStateBackend.createKeyedStateBackend(EmbeddedRocksDBStateBackend.java:484)

    at org.apache.flink.contrib.streaming.state.EmbeddedRocksDBStateBackend.createKeyedStateBackend(EmbeddedRocksDBStateBackend.java:97)

    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: java.io.IOException: Could not find class 'old.class.path' in classpath.

    at org.apache.flink.util.InstantiationUtil.resolveClassByName(InstantiationUtil.java:775)

    at org.apache.flink.util.InstantiationUtil.resolveClassByName(InstantiationUtil.java:750)

    at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializerSnapshotData.readTypeClass(KryoSerializerSnapshotData.java:186)

    at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializerSnapshotData.createFrom(KryoSerializerSnapshotData.java:69)

    at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializerSnapshot.readSnapshot(KryoSerializerSnapshot.java:81)

    at org.apache.flink.api.common.typeutils.TypeSerializerSnapshot.readVersionedSnapshot(TypeSerializerSnapshot.java:175)

    at org.apache.flink.api.common.typeutils.NestedSerializersSnapshotDelegate.readNestedSerializerSnapshots(NestedSerializersSnapshotDelegate.java:178)

    at org.apache.flink.api.common.typeutils.CompositeTypeSerializerSnapshot.readSnapshot(CompositeTypeSerializerSnapshot.java:171)

    at org.apache.flink.api.common.typeutils.TypeSerializerSnapshot.readVersionedSnapshot(TypeSerializerSnapshot.java:175)

    at org.apache.flink.api.common.typeutils.TypeSerializerSnapshotSerializationUtil$TypeSerializerSnapshotSerializationProxy.deserializeV2(TypeSerializerSnapshotSerializationUtil.java:174)

    at org.apache.flink.api.common.typeutils.TypeSerializerSnapshotSerializationUtil$TypeSerializerSnapshotSerializationProxy.read(TypeSerializerSnapshotSerializationUtil.java:145)

    at org.apache.flink.api.common.typeutils.TypeSerializerSnapshotSerializationUtil.readSerializerSnapshot(TypeSerializerSnapshotSerializationUtil.java:77)

    at org.apache.flink.runtime.state.metainfo.StateMetaInfoSnapshotReadersWriters$CurrentReaderImpl.readStateMetaInfoSnapshot(StateMetaInfoSnapshotReadersWriters.java:237)

    at org.apache.flink.runtime.state.KeyedBackendSerializationProxy.read(KeyedBackendSerializationProxy.java:184)

    at org.apache.flink.runtime.state.restore.FullSnapshotRestoreOperation.readMetaData(FullSnapshotRestoreOperation.java:194)

    at org.apache.flink.runtime.state.restore.FullSnapshotRestoreOperation.restoreKeyGroupsInStateHandle(FullSnapshotRestoreOperation.java:171)

    at org.apache.flink.runtime.state.restore.FullSnapshotRestoreOperation.access$100(FullSnapshotRestoreOperation.java:113)

    at org.apache.flink.runtime.state.restore.FullSnapshotRestoreOperation$1.next(FullSnapshotRestoreOperation.java:158)

    at org.apache.flink.runtime.state.restore.FullSnapshotRestoreOperation$1.next(FullSnapshotRestoreOperation.java:140)

    at org.apache.flink.contrib.streaming.state.restore.RocksDBFullRestoreOperation.restore(RocksDBFullRestoreOperation.java:102)

    at org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackendBuilder.build(RocksDBKeyedStateBackendBuilder.java:315)

    ... 18 more

Caused by: java.lang.ClassNotFoundException: old.class.path

    at java.base/java.net.URLClassLoader.findClass(URLClassLoader.java:476)

    at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:589)

    at org.apache.flink.util.FlinkUserCodeClassLoader.loadClassWithoutExceptionHandling(FlinkUserCodeClassLoader.java:68)

    at org.apache.flink.util.ChildFirstClassLoader.loadClassWithoutExceptionHandling(ChildFirstClassLoader.java:74)

    at org.apache.flink.util.FlinkUserCodeClassLoader.loadClass(FlinkUserCodeClassLoader.java:52)

    at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:522)

    at org.apache.flink.runtime.execution.librarycache.FlinkUserCodeClassLoaders$SafetyNetWrapperClassLoader.loadClass(FlinkUserCodeClassLoaders.java:177)

    at java.base/java.lang.Class.forName0(Native Method)

    at java.base/java.lang.Class.forName(Class.java:398)

    at org.apache.flink.util.InstantiationUtil.resolveClassByName(InstantiationUtil.java:773)

    ... 38 more

迁移后的Operator代码

import new.class.path.MyClass;
//import old.new.class.path.MyClass; // old class path

import org.apache.flink.api.common.state.StateTtlConfig;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;

public class Operator extends KeyedProcessFunction<String, OperatorInput, OperatorOutput> {
    private ValueState<MyClass> myClassState;
    public void open(Configuration parameters) {
        ValueStateDescriptor<MyClass> myClassCacheDescriptor = new ValueStateDescriptor<>("State", MyClass.class);
        myClassCacheDescriptor.enableTimeToLive(ttlConfig);
        this.myClassState= getRuntimeContext().getState(myClassCacheDescriptor);

    }
}

官方文档说明

Flink 1.15官方文档明确说明:
POJO类型的类名(包括类的命名空间)不能更改。


解决方案

方案一:添加类路径映射(推荐)

在Flink配置中添加类重定向规则,让Flink在反序列化时将旧类路径映射到新路径。在AWS KDA中,可通过应用配置的flink-conf.yaml添加以下配置:

classloader.resolve-order: parent-first
classloader.mappings: old.class.path.MyClass=new.class.path.MyClass

该配置会告知Flink类加载器,当尝试加载old.class.path.MyClass时,实际使用new.class.path.MyClass替代,无需修改业务代码。

方案二:保留旧类占位符

在新代码中保留旧路径的空类作为占位符,并让它继承新类,确保序列化兼容:

// 放在old.class.path包下的占位类
package old.class.path;
import new.class.path.MyClass;

public class MyClass extends new.class.path.MyClass {}

这样Flink恢复时能找到旧类路径的类,实际逻辑复用新类实现,避免类找不到的异常。

方案三:自定义序列化器

为MyClass实现自定义TypeSerializer,在序列化和反序列化时手动处理类路径映射。需在ValueStateDescriptor中指定自定义序列化器:

ValueStateDescriptor<MyClass> myClassCacheDescriptor = new ValueStateDescriptor<>(
    "State",
    MyClass.class,
    new MyCustomSerializer() // 自定义序列化器
);

自定义序列化器需实现TypeSerializer接口,在反序列化时识别旧类路径的数据并转换为新类实例,适合复杂类结构变更场景。


  • 方案一无需修改代码,是AWS KDA托管环境下最简洁的解决方式;
  • 方案二需维护占位类,适合临时过渡;
  • 方案三灵活性最高,但实现成本较高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:49:53