Scala 2.12中Priority Queue Kryo序列化器创建失败问题
Scala 2.12下Kryo序列化Priority Queue时的InstantiationException问题
问题现象
在Scala 2.12环境中自定义mutable.PriorityQueue的Kryo序列化器时,触发InstantiationException异常,尝试通过对象实例化方式解决后问题依旧。
报错堆栈
Caused by: java.lang.InstantiationException: PriorityQueueSerializer at java.base/java.lang.Class.newInstance(Class.java:671) at com.esotericsoftware.kryo.factories.ReflectionSerializerFactory.makeSerializer(ReflectionSerializerFactory.java:43) ... 63 more Caused by: java.lang.NoSuchMethodException: PriorityQueueSerializer.<init>() at java.base/java.lang.Class.getConstructor0(Class.java:3617) at java.base/java.lang.Class.newInstance(Class.java:658) ... 64 more
自定义序列化器代码
class PriorityQueueSerializer[T](implicit ord: Ordering[T]) extends Serializer[mutable.PriorityQueue[T]](true, false) { override def write(kryo: Kryo, output: Output, pq: mutable.PriorityQueue[T]): Unit = { println("Write is running") output.writeInt(pq.size) pq.foreach { elem => kryo.writeClassAndObject(output, elem) } kryo.writeClassAndObject(output, ord) } override def read(kryo: Kryo, input: Input, clazz: Class[mutable.PriorityQueue[T]]): mutable.PriorityQueue[T] = { println("Read is running") val pq = new mutable.PriorityQueue[T]()(ord) kryo.reference(pq) input.readStringBuilder() val size = input.readInt() for (_ <- 0 until size) { val elem = kryo.readClassAndObject(input).asInstanceOf[T] pq.enqueue(elem) } val ordering = kryo.readClassAndObject(input).asInstanceOf[Ordering[T]] new mutable.PriorityQueue()(ordering) ++= pq } }
背景原因
升级Scala版本从2.11到2.12后,在Flink中遇到空指针异常,因此尝试自定义Kryo序列化器解决该问题。
解决方案
问题根源
报错明确指出PriorityQueueSerializer缺少无参构造方法,Kryo的ReflectionSerializerFactory默认会通过无参构造器实例化序列化器,但当前序列化器的构造方法依赖隐式Ordering[T],无法被Kryo直接实例化。此外原代码中input.readStringBuilder()是多余操作,会导致读数据时位置偏移,引发后续解析错误。
修复方案
方案一:调整序列化器逻辑,移除构造依赖
修改序列化器,直接使用PriorityQueue自带的Ordering进行读写,同时删除多余的读操作:
class PriorityQueueSerializer[T] extends Serializer[mutable.PriorityQueue[T]](true, false) { override def write(kryo: Kryo, output: Output, pq: mutable.PriorityQueue[T]): Unit = { // 先写入队列大小,再写入元素,最后写入队列自带的排序规则 output.writeInt(pq.size) pq.foreach(elem => kryo.writeClassAndObject(output, elem)) kryo.writeClassAndObject(output, pq.ordering) } override def read(kryo: Kryo, input: Input, clazz: Class[mutable.PriorityQueue[T]]): mutable.PriorityQueue[T] = { val size = input.readInt() // 先读取排序规则,再初始化队列 val ordering = kryo.readClassAndObject(input).asInstanceOf[Ordering[T]] val pq = mutable.PriorityQueue.empty[T](ordering) kryo.reference(pq) // 批量读取元素并加入队列 for (_ <- 0 until size) { val elem = kryo.readClassAndObject(input).asInstanceOf[T] pq.enqueue(elem) } pq } }
这种方式让序列化器拥有无参构造,同时保证读写逻辑的一致性,无需依赖外部隐式参数。
方案二:注册序列化器时传入实例
如果需要保留原构造逻辑,注册序列化器时不要传入Class,而是直接传入序列化器实例:
// 以String类型为例,需为每个泛型类型单独注册 kryo.register(classOf[mutable.PriorityQueue[String]], new PriorityQueueSerializer[String])
此方式适合泛型类型固定的场景,灵活性较低。
内容的提问来源于stack exchange,提问作者ashwin gupta
相关产品推荐
相关产品推荐

