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

Flink 1.12中Java DataSet与Scala DataStream的Key序列化器不兼容问题

问题背景

使用Flink 1.12版本的Scala代码开发流处理程序,包含按键时间窗口、处理函数逻辑,依赖外部化保存点实现故障恢复。当采用仅支持Java的DataSet API状态处理器进行状态操作与迁移时,出现序列化兼容性问题:流处理程序使用ScalaCaseClassSerializer序列化Scala Case Class类型的Key,而DataSet API默认使用KryoSerializer读取,导致报错:

org.apache.flink.util.StateMigrationException: The new key serializer (org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer@4541edf4) must be compatible with the previous key serializer (org.apache.flink.api.scala.typeutils.ScalaCaseClassSerializer@d1821229).

已知切换到Java DataStream API或升级Flink至1.15+版本可解决问题,但前者不符合当前技术选型,后者会带来大量代码重构,且enableForceKryo()配置无效,需在Flink 1.12版本内找到解决方案。


可行解决方案

方案1:为DataSet环境手动注册Scala Case Class序列化器

在状态读取程序中,强制ExecutionEnvironment对Scala Case Class类型使用ScalaCaseClassSerializer,而非默认的Kryo序列化器,确保与流处理程序的序列化逻辑一致。

代码示例:

import org.apache.flink.api.scala.typeutils.ScalaCaseClassSerializer
import org.apache.flink.api.common.typeinfo.TypeInformation
import org.apache.flink.runtime.state.StateBackend
import org.apache.flink.runtime.state.memory.MemoryStateBackend
import org.apache.flink.api.java.operators.Savepoint

def main(args: Array[String]): Unit = {
  val checkpointDir = "/path/to/above/checkpoint/dir"
  val env = ExecutionEnvironment.getExecutionEnvironment
  
  // 为Scala Case Class注册对应序列化器
  val keyTypeInfo = TypeInformation.of(classOf[Key])
  val keySerializer = new ScalaCaseClassSerializer[Key](keyTypeInfo.getTypeClass, env.getConfig)
  env.getConfig.registerTypeWithKryoSerializer(classOf[Key], keySerializer.getClass)
  
  val backend: StateBackend = new MemoryStateBackend()
  val savepoint = Savepoint.load(env, checkpointDir, backend)

  savepoint
    .readKeyedState("aggregate", new StateReaderFunction)
    .printToErr()

  env.execute()
}

方案2:自定义兼容的Key序列化器并双向配置

定义完全可控的自定义序列化器,同时在流处理程序和状态读取程序中配置使用,彻底统一两端的序列化逻辑。

  1. 自定义序列化器:
import org.apache.flink.api.common.typeutils.TypeSerializer
import org.apache.flink.core.memory.{DataInputView, DataOutputView}

class KeyCustomSerializer extends TypeSerializer[Key] {
  override def isImmutableType: Boolean = true
  override def duplicate(): TypeSerializer[Key] = this
  override def serialize(record: Key, target: DataOutputView): Unit = {
    target.writeUTF(record.prefix)
    target.writeUTF(record.suffix)
  }
  override def deserialize(source: DataInputView): Key = {
    Key(source.readUTF(), source.readUTF())
  }
  override def deserialize(reuse: Key, source: DataInputView): Key = deserialize(source)
  override def copy(from: Key): Key = from
  override def copy(from: Key, reuse: Key): Key = from
  override def getLength: Int = -1
}
  1. 在流处理程序中注册:
// 创建StreamExecutionEnvironment后添加
env.getConfig.registerTypeWithKryoSerializer(classOf[Key], classOf[KeyCustomSerializer])
  1. 在状态读取程序中注册:
// 创建ExecutionEnvironment后添加
env.getConfig.registerTypeWithKryoSerializer(classOf[Key], classOf[KeyCustomSerializer])

注意:修改流处理程序的序列化器后,需重新生成保存点;若要兼容旧保存点,需确保自定义序列化器能正确反序列化旧数据。

方案3:使用Scala版DataSet ExecutionEnvironment

Flink 1.12的Scala API对DataSet有原生支持,使用Scala版ExecutionEnvironment会自动优先使用Scala相关类型序列化器(包括ScalaCaseClassSerializer),与流处理程序的序列化逻辑对齐。

代码示例:

import org.apache.flink.api.scala.ExecutionEnvironment
import org.apache.flink.api.scala.typeutils.Types
import org.apache.flink.runtime.state.StateBackend
import org.apache.flink.runtime.state.memory.MemoryStateBackend
import org.apache.flink.api.java.operators.Savepoint

def main(args: Array[String]): Unit = {
  val checkpointDir = "/path/to/above/checkpoint/dir"
  // 使用Scala版ExecutionEnvironment
  val env = ExecutionEnvironment.getExecutionEnvironment
  val backend: StateBackend = new MemoryStateBackend()

  val savepoint = Savepoint.load(env, checkpointDir, backend)
  
  // 显式指定Key类型信息,确保使用Scala序列化器
  savepoint
    .readKeyedState("aggregate", new StateReaderFunction, Types.of[Key])
    .printToErr()

  env.execute()
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 08:08:10