Spark 3.5.0中mapPartitions内可变列表迭代后为空是否正常?
问题解答:mapPartitions中ListBuffer为空的原因
这不是预期现象,问题核心在于Scala Iterator的惰性求值特性:
- 你代码里的
it.map(row => { list += row; row })返回的是一个未被遍历的惰性Iterator,map中的逻辑(往ListBuffer添加元素)只有在这个Iterator被实际迭代时才会执行。 - 你在创建完
result后立刻打印list.size,此时result还没被Spark触发遍历,map里的代码根本没运行,所以ListBuffer是空的。 - 后续调用
.show()时,Spark才会真正迭代result,此时ListBuffer才会被填充,但打印语句已经执行完毕,所以你看不到正确的size。
验证与修改方案
如果要在打印时看到正确的ListBuffer大小,可以主动触发Iterator的遍历,比如将result先转为List再转回Iterator:
import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder import scala.collection.mutable val encoder = ExpressionEncoder(df.schema) df.mapPartitions(it => { val list = mutable.ListBuffer[Row]() // 先遍历Iterator,触发map逻辑填充ListBuffer val resultList = it.map(row => { list += row row }).toList println(s"list.size: ${list.size}") // 转回Iterator返回给Spark resultList.iterator })(encoder).show()
或者直接用foreach遍历原Iterator来填充ListBuffer:
import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder import scala.collection.mutable val encoder = ExpressionEncoder(df.schema) df.mapPartitions(it => { val list = mutable.ListBuffer[Row]() it.foreach(row => list += row) println(s"list.size: ${list.size}") list.iterator })(encoder).show()
内容的提问来源于stack exchange,提问作者Antonio Ye
相关产品推荐
相关产品推荐

