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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:22:33