Spark中rdd.foreachPartition内部的foreach()不生效是什么原因
问题原因
Scala 中的 Iterator 是仅支持单次遍历的结构,它本质是指向集合元素的游标,本身不持有任何元素数据:
- 当你调用
records.length时,程序会一次性遍历整个迭代器的所有元素完成计数,遍历结束后游标已经走到迭代器末尾 - 后续执行
records.foreach时,迭代器已经没有剩余元素可以读取,自然不会输出任何内容
这个问题是 Scala 语言的 Iterator 特性导致的,和 Spark 框架本身的逻辑无关,不是你对 foreachPartition 传参规则的理解有误。
解决方案
如果你需要同时获取分区数据长度和遍历处理元素,有两种可选方案:
方案1:小数据量场景,转成可重复访问的集合
如果单分区数据量不大,可以先把迭代器转成 List 等常驻内存的集合,就能支持多次读取:
stream .map(x => x.value()) .foreachRDD( rdd => { rdd.foreachPartition( (records: Iterator[String]) => { val recordList = records.toList println(recordList.length) recordList.foreach(x => println(x)) } ) } )
方案2:大数据量场景,单次遍历完成计数+处理
如果单分区数据量很大,直接转 List 可能导致 Executor 内存溢出,可以在一次遍历里同时完成计数和处理逻辑,不额外占用内存:
stream .map(x => x.value()) .foreachRDD( rdd => { rdd.foreachPartition( (records: Iterator[String]) => { var count = 0 records.foreach(x => { count += 1 println(x) }) println(count) } ) } )
内容的提问来源于stack exchange,提问作者cnidaye
相关产品推荐
相关产品推荐

