升级Flink至1.16.2后遇FlinkScalaKryoInstantiator类找不到异常求助
Flink 1.16.2升级后出现ClassNotFoundException: org.apache.flink.runtime.types.FlinkScalaKryoInstantiator
我在使用Apache Flink进行流数据处理,之前用的是Flink 1.13.16,升级到1.16.2后触发如下异常:
12:38:38.845 [main] INFO org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer - Kryo serializer scala extensions are not available. java.lang.ClassNotFoundException: org.apache.flink.runtime.types.FlinkScalaKryoInstantiator at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:581) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:178) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:522) at java.base/java.lang.Class.forName0(Native Method) at java.base/java.lang.Class.forName(Class.java:315) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.getKryoInstance(KryoSerializer.java:486) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.checkKryoInitialized(KryoSerializer.java:521) at org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer.copy(KryoSerializer.java:306) at org.apache.flink.api.java.typeutils.runtime.PojoSerializer.copy(PojoSerializer.java:250) at org.apache.flink.api.java.typeutils.runtime.PojoSerializer.copy(PojoSerializer.java:250) at org.apache.flink.streaming.util.AbstractStreamOperatorTestHarness$MockOutput.collect(AbstractStreamOperatorTestHarness.java:846) at org.apache.flink.streaming.util.AbstractStreamOperatorTestHarness$MockOutput.collect(AbstractStreamOperatorTestHarness.java:807) at org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:56) at org.apache.flink.streaming.api.operators.CountingOutput.collect(CountingOutput.java:29) at org.apache.flink.streaming.api.operators.TimestampedCollector.collect(TimestampedCollector.java:51) at com.abp.accountstatements.projections.PcBasedOnCalendarChange.processElement(PcBasedOnCalendarChange.java:59)
相关代码
PcBasedOnCalendarChange.java
public class PcBasedOnCalendarChange extends KeyedBroadcastProcessFunction< String, GenerationDetail, String, GenerationDetail> { ..... ..... @Override public void processElement( GenerationDetail genDetail, ReadOnlyContext ctx, Collector<GenerationDetail> collector) throws Exception { try { String accNum = genDetail.getAccount().getAccountNumber(); statementGenDateState.update(genDetail); collector.collect(genDetail); // 升级至Flink 1.16.2后此处抛出上述异常 } catch (Exception e) { log.error(e.getMessage()); ctx.output(OutputExceptions.statementGenerationDetailOutput, genDetail); } } ..... }
GenerationDetail.java
@Data @JsonIgnoreProperties(ignoreUnknown = true) public class GenerationDetail { private String accountId; private AccountIdentifier account; // AccountIdentifier是可序列化类 private BusinessDateDTO accOpeningDate; // BusinessDateDTO是可序列化类 private String generateFrequency; private String contractId; private String productId; private String calendarCode; private Frequency frequency; // Frequency是可序列化类 }
已知信息与尝试过的操作
org.apache.flink.runtime.type包从Flink 1.14.0版本起已被移除- 尝试过设置不同的序列化器,但问题未解决
- 已参考过Flink序列化调优相关文档
解决方案
原因分析
Flink 1.14+彻底移除了FlinkScalaKryoInstantiator类,但你的环境中可能存在以下情况导致Kryo仍尝试加载该类:
- 依赖中混合了不同版本的Flink组件,旧版本的Scala相关依赖残留
- Kryo配置中显式指定了该已被移除的实例化类
- 测试代码使用的
AbstractStreamOperatorTestHarness依赖版本不一致,触发了旧的序列化逻辑
具体解决步骤
统一Flink依赖版本
确保项目中所有Flink相关依赖(包括核心、流处理、测试组件)的版本均为1.16.2,通过依赖管理工具(如Maven的dependencyManagement)锁定版本,排除传递依赖中引入的旧版本Flink包。清理Kryo配置
检查代码或配置文件中是否有指定FlinkScalaKryoInstantiator的配置,若有则移除。Flink 1.16+默认的Kryo实例化方式已不再依赖该类。显式禁用Kryo Scala扩展
在代码中添加配置,强制关闭Kryo的Scala扩展支持:Configuration config = new Configuration(); config.setString("kryo.scala.extensions.enabled", "false"); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);确保Pojo序列化正常
- 确认
GenerationDetail及其内部类(AccountIdentifier、BusinessDateDTO、Frequency)均符合Flink Pojo序列化要求:类是public、有默认无参构造、所有字段public或有getter/setter - 若仍触发Kryo序列化,可显式为这些类注册Pojo序列化器:
env.getConfig().registerPojoType(GenerationDetail.class); env.getConfig().registerPojoType(AccountIdentifier.class); // 其他内部类同理
- 确认
更新测试依赖
确保测试代码中使用的flink-streaming-java-test等测试依赖版本为1.16.2,避免旧版本测试 harness引入的序列化逻辑冲突。
内容的提问来源于stack exchange,提问作者Ajit Gupta
相关产品推荐
相关产品推荐

