Spark中遍历flatMap及在for循环内执行flatMap的技术问询
Spark中flatMap结果遍历与循环内使用的实现指导
一、如何遍历Spark中flatMap的结果?
首先得明确:Spark的分布式数据集(不管是DStream、RDD还是DataFrame)都不能像本地List那样直接用for循环遍历——因为数据分散在集群各个节点上,不是都在Driver端。针对不同场景,有两种靠谱的处理方式:
- 拉取到本地遍历(小数据量适用)
如果flatMap后的数据集不大,可以用行动算子collect()把数据拉到Driver端,然后在本地遍历处理。注意:数据量大的话别这么干,会把Driver内存撑爆。示例代码:
val flatMappedStream = yourOriginalDStream.flatMap(item => someProcessing(item)) flatMappedStream.foreachRDD { rdd => // 先把当前批次的RDD数据拉到Driver本地 val localData = rdd.collect() localData.foreach(item => { // 这里就能像本地集合一样遍历处理了,比如打印、写入本地文件等 println(s"处理元素: $item") }) }
- 分布式遍历处理(大数据量适用)
如果数据量很大,不想拉到本地,就用foreach算子在Executor节点上直接处理每个元素,全程分布式执行:
flatMappedStream.foreachRDD { rdd => rdd.foreach(item => { // 每个Executor节点上独立处理元素,比如写入数据库、分布式计算等 writeToDatabase(item) }) }
⚠️ 注意:DStream是流式数据,所有操作都要基于foreachRDD来处理每个批次的RDD,不能直接遍历DStream对象本身哦。
二、在Spark for循环中使用flatMap的正确姿势
先看你写的代码,这里面有几个容易踩的坑:
- 在
flatMap里修改Driver端的best和bestVarianceSum变量完全无效——flatMap的逻辑是在Executor节点跑的,Driver的变量副本不会被同步更新,而且多个Executor同时改还会有并发冲突。 flatMap是转换算子,不会触发实际计算,你的循环里只是定义了转换逻辑,没有行动算子的话,代码根本不会执行。- 你想要的是多次聚类后选出方差最小的最优结果,这个逻辑不能在分布式的
flatMap里操作Driver变量实现,得换个思路。
修正后的实现代码
我们把聚类后的结果拉到Driver端做比较(如果数据量不大),或者把方差计算分布式执行后只拉取结果到Driver比较,下面给两种场景的示例:
场景1:每个批次数据量不大,拉到Driver处理
import org.apache.spark.streaming.dstream.DStream import scala.collection.mutable.ArrayBuffer def clusterVar(points: DStream[ArrayBuffer[Double]]): DStream[ArrayBuffer[Double]]={ // 用transform算子,允许我们在Driver端操作每个批次的RDD points.transform { rdd => if (rdd.isEmpty()) { // 空批次直接返回空RDD spark.sparkContext.emptyRDD[ArrayBuffer[Double]] } else { // 把当前批次的点数据拉到Driver本地 val localPoints = rdd.collect() var bestCluster: ArrayBuffer[Double] = ArrayBuffer.empty var bestVarianceSum = Double.PositiveInfinity // 多次尝试聚类,在Driver端比较最优结果 for (i <- 0 until numTrials) { val currentCluster = cluster(localPoints) // 这里cluster函数处理本地数据 val currentVariance = score(currentCluster) // 计算当前聚类的方差和 if (currentVariance < bestVarianceSum) { bestVarianceSum = currentVariance bestCluster = currentCluster.clone() // 克隆避免引用传递的问题 } } // 把最优聚类结果转为RDD,作为当前批次的输出 spark.sparkContext.parallelize(Seq(bestCluster)) } } }
场景2:数据量很大,分布式聚类+Driver端比较方差
如果每个批次的数据量太大,不能拉到Driver,就把聚类和方差计算放在集群上执行,只把方差结果拉到Driver比较:
def clusterVar(points: DStream[ArrayBuffer[Double]]): DStream[ArrayBuffer[Double]]={ points.transform { rdd => if (rdd.isEmpty()) { spark.sparkContext.emptyRDD[ArrayBuffer[Double]] } else { var bestVarianceSum = Double.PositiveInfinity var bestClusterRDD = spark.sparkContext.emptyRDD[ArrayBuffer[Double]] for (i <- 0 until numTrials) { val clusterRDD = cluster(rdd) // cluster函数返回分布式的RDD[ArrayBuffer[Double]] // 分布式计算方差和:每个聚类计算方差,再求和 val currentVarianceSum = clusterRDD.map(cluster => score(cluster)).reduce(_ + _) if (currentVarianceSum < bestVarianceSum) { bestVarianceSum = currentVarianceSum bestClusterRDD = clusterRDD } } bestClusterRDD } } }
关键要点总结
- 永远不要在Executor端的算子(比如
flatMap、map)里修改Driver端的变量,分布式环境下这种操作既无效又容易出问题。 - 用
transform算子可以让你在Driver端操作每个批次的RDD,非常适合这种需要多次尝试后选最优的场景。 - 最终的DStream需要通过行动算子(比如
print()、saveAsTextFiles())触发实际计算,不然所有转换逻辑都是“纸上谈兵”。
内容的提问来源于stack exchange,提问作者Suzy Tros
相关产品推荐
相关产品推荐

