如何在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
相关产品推荐
相关产品推荐

