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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 23:33:18