Scala Spark mapGroup报错:No Encoder found for Iterator[String]求助
问题解决:Spark Dataset合并时Iterator编码错误
问题背景
现有Dept和Sal两个Dataset,需要合并成包含部门信息、下属薪资列表及薪资总和的目标表。执行joinWith+groupByKey后,尝试转换为自定义case class时触发以下错误:
java.lang.UnsupportedOperationException: No Encoder found for Iterator[String] - field (class: "scala.collection.Iterator", name: "_2") - root class: "scala.Tuple3"
现有代码与数据集
定义的case class及数据集初始化代码:
case class DeptSals(deptNo :Int, dname: String ,sals : Seq[Sal]) case class Dept(deptNo :Int, dname: String) case class Sal(deptNo :Int, Ename: String,sal:Long) val deptDS = Seq(Dept(1,"ANALYST"),Dept(2,"HR"),Dept(3,"FINANCE"),Dept(4,"CLERK"),Dept(5,"MANAGER")).toDS val salDS = Seq(Sal(2,"Amelia",4000),Sal(2,"Cherry",4000),Sal(3,"Carl",7000),Sal(4,"John",6000),Sal(5,"Rick",1000)).toDS
Dept数据集:
scala> deptDS.show +------+-------+ |deptNo| dname| +------+-------+ | 1|ANALYST| | 2| HR| | 3|FINANCE| | 4| CLERK| | 5|MANAGER| +------+-------+
Sal数据集:
+------+------+----+ |deptNo| Ename| sal| +------+------+----+ | 2|Cherry|4000| | 2|Amelia|4000| | 3| Carl|7000| | 4| John|6000| | 5| Rick|1000| +------+------+----+
出错的执行代码:
import spark.implicits._ val deptJoinSalDS = deptDS.joinWith(salDS,deptDS("deptNo") === salDS("deptNo"),"left_outer") .groupByKey(x => x._1.deptNo) .mapGroups{ case(k,v) => ( k,v.map( _._1.dname), v.map( _._2.Ename)) }
期望输出
+------+---------------------------------+-------+ |deptNo|Sals | sumSal| +------+---------------------------------+-------+ | 1| | 0| | 2|[[2,Amelia,4000],[2,Cherry,4000]]| 8000| | 3|[[3,Carl,7000]] | 7000| | 4|[[4,John,6000]] | 6000| | 5|[[5,Rick,1000]] | 1000| +------+---------------------------------+-------+
错误原因
mapGroups中的参数v是Iterator类型,Spark没有为Iterator提供内置Encoder:
- Iterator是惰性遍历的集合,只能被消费一次,Spark无法对其进行序列化、持久化或重复遍历
- Spark的Encoder只支持可序列化、可重复访问的集合类型(如
Seq、List)
解决方案
1. 定义匹配期望输出的case class
首先调整case class结构,对齐期望的三个字段:
case class DeptSalSummary(deptNo: Int, sals: Seq[Sal], sumSal: Long)
2. 修改mapGroups逻辑
将Iterator转换为Seq,同时处理left join后无薪资数据的情况:
import spark.implicits._ val resultDS = deptDS.joinWith(salDS, deptDS("deptNo") === salDS("deptNo"), "left_outer") .groupByKey(_._1.deptNo) .mapGroups { case (deptNo, iter) => // 将Iterator转为Seq,方便多次访问 val groupedData = iter.toSeq // 提取薪资列表,过滤掉left join产生的null值 val salList = groupedData.flatMap(_._2) // 计算薪资总和,无数据时返回0 val totalSal = salList.map(_.sal).sum // 返回自定义case class实例 DeptSalSummary(deptNo, salList, totalSal) } // 查看结果 resultDS.show(false)
3. 执行结果
运行后将得到符合期望的输出:
+------+---------------------------------+-------+ |deptNo|sals |sumSal | +------+---------------------------------+-------+ |1 |[] |0 | |2 |[Sal(2,Amelia,4000), Sal(2,Cherry,4000)]|8000| |3 |[Sal(3,Carl,7000)] |7000 | |4 |[Sal(4,John,6000)] |6000 | |5 |[Sal(5,Rick,1000)] |1000 | +------+---------------------------------+-------+
额外优化:使用agg替代mapGroups(更高效)
对于这类聚合场景,使用Spark内置的聚合函数性能更优,无需手动处理Iterator:
import org.apache.spark.sql.functions._ val resultDF = deptDS.join(salDS, Seq("deptNo"), "left_outer") .groupBy("deptNo", "dname") .agg( collect_list(struct("deptNo", "Ename", "sal")).alias("sals"), coalesce(sum("sal"), lit(0)).alias("sumSal") ) // 转换为Dataset(可选) resultDF.as[DeptSalSummary].show(false)
内容的提问来源于stack exchange,提问作者Andy
相关产品推荐
相关产品推荐

