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

升级Flink至1.16.2后遇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仍尝试加载该类:

  1. 依赖中混合了不同版本的Flink组件,旧版本的Scala相关依赖残留
  2. Kryo配置中显式指定了该已被移除的实例化类
  3. 测试代码使用的AbstractStreamOperatorTestHarness依赖版本不一致,触发了旧的序列化逻辑

具体解决步骤

  1. 统一Flink依赖版本
    确保项目中所有Flink相关依赖(包括核心、流处理、测试组件)的版本均为1.16.2,通过依赖管理工具(如Maven的dependencyManagement)锁定版本,排除传递依赖中引入的旧版本Flink包。

  2. 清理Kryo配置
    检查代码或配置文件中是否有指定FlinkScalaKryoInstantiator的配置,若有则移除。Flink 1.16+默认的Kryo实例化方式已不再依赖该类。

  3. 显式禁用Kryo Scala扩展
    在代码中添加配置,强制关闭Kryo的Scala扩展支持:

    Configuration config = new Configuration();
    config.setString("kryo.scala.extensions.enabled", "false");
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
    
  4. 确保Pojo序列化正常

    • 确认GenerationDetail及其内部类(AccountIdentifier、BusinessDateDTO、Frequency)均符合Flink Pojo序列化要求:类是public、有默认无参构造、所有字段public或有getter/setter
    • 若仍触发Kryo序列化,可显式为这些类注册Pojo序列化器:
      env.getConfig().registerPojoType(GenerationDetail.class);
      env.getConfig().registerPojoType(AccountIdentifier.class);
      // 其他内部类同理
      
  5. 更新测试依赖
    确保测试代码中使用的flink-streaming-java-test等测试依赖版本为1.16.2,避免旧版本测试 harness引入的序列化逻辑冲突。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 08:24:56