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

如何在Spark Aggregator中使用嵌套Map作为缓冲区并解决编码器问题

Scala Spark Aggregator中Map[String, Set[String]]编码器问题

我在Scala Spark中实现Aggregator时,希望使用Map[String, Set[String]]作为缓冲区类型。单独对Set[String]使用kryo或ExpressionEncoder编码是可行的,但将其嵌入Map后,系统无法找到对应的编码器。

我尝试了以下两种编码器定义方式:

def bufferEncoder: Encoder[Map[String, Set[String]]] = Encoders.kryo[Map[String, Set[String]]]
def bufferEncoder: Encoder[Map[String, Set[String]]] = implicitly(ExpressionEncoder[Map[String, Set[String]]])

此前编写的另一个Aggregator使用Encoders.kryo[Set[String]]作为缓冲区编码器可以正常运行,但上述两种写法均触发错误:

java.lang.UnsupportedOperationException: No Encoder found for java.util.Map[String,Array[String]]

更新内容

我添加了两个极简代码示例,唯一区别在于一个使用Map[String, Array[String]],另一个使用Map[String, Set[String]]。前者可编译运行(结果无关紧要),后者则触发上述异常。请注意,两个示例均使用可变Scala Map和Set。

MapSetTest代码

class MapSetTest extends Aggregator[String, Map[String, Set[String]], Int] with Serializable {

override def zero = Map[String, Set[String]]()

override def reduce(buffer: Map[String, Set[String]], newItem: String) = {
    buffer.put(newItem, Set[String]() + newItem)
    buffer
}


override def merge(b1: Map[String, Set[String]], b2: Map[String, Set[String]]) = {
  b1
}

override def finish(reduction: Map[String, Set[String]]): Int = {
  reduction.size
}

def bufferEncoder = implicitly[Encoder[Map[String, Set[String]]]]
def outputEncoder = Encoders.scalaInt
}

MapArrayTest代码

class MapArrayTest extends Aggregator[String, Map[String, Array[String]], Int] with Serializable {

override def zero = Map[String, Array[String]]()

override def reduce(buffer: Map[String, Array[String]], newItem: String) = {
  buffer.put(newItem, Array[String](newItem))
  buffer
}


override def merge(b1: Map[String, Array[String]], b2: Map[String, Array[String]]) = {
  b1
}

override def finish(reduction: Map[String, Array[String]]): Int = {
  reduction.size
}

def bufferEncoder = implicitly[Encoder[Map[String, Array[String]]]]

def outputEncoder = Encoders.scalaInt
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 18:46:06