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序列化器并双向配置
定义完全可控的自定义序列化器,同时在流处理程序和状态读取程序中配置使用,彻底统一两端的序列化逻辑。
- 自定义序列化器:
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 }
- 在流处理程序中注册:
// 创建StreamExecutionEnvironment后添加 env.getConfig.registerTypeWithKryoSerializer(classOf[Key], classOf[KeyCustomSerializer])
- 在状态读取程序中注册:
// 创建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

