Spark中RDD内遍历广播变量时如何获取正确的RDD map输出
问题根因
- 代码中
map算子内部的for循环未携带yield关键字:Scala语法中不带yield的for推导式默认返回Unit类型,序列化后写入文本就表现为空括号(),这是问题的直接诱因。 - 即使补充
yield,map算子也不符合需求:map是1对1转换算子,每个输入元素仅会返回1个输出值,补充yield后会得到每个输入对应一个元组集合,最终输出为嵌套结构,和你预期的扁平化独立元组不符,需要改用flatMap(1对多转换算子,会自动把返回的集合拆分为单个元素输出)。
另外你原演示代码还存在两处语法错误:rdd1的parallelize方法调用缺右括号,rdd2的类型声明前多了等于号,修复代码时会一并处理。
修复后代码
import org.apache.spark.SparkConf import org.apache.spark.SparkContext import org.apache.spark.rdd.RDD val conf = new SparkConf().setAppName("Sample").setMaster("local") val sc = new SparkContext(conf) val rdd1: RDD[(Int, Int)] = sc.parallelize(Seq((1,1),(2,3),(3,4))) val rdd2: RDD[(Int, Int)] = sc.parallelize(Seq((2,1),(3,2),(7,6))) val rdd2AsMap = rdd2.collectAsMap.toMap val broadcastMap = sc.broadcast(rdd2AsMap) // 替换map为flatMap,for循环添加yield val result = rdd1.flatMap{case(x,y) => { for((key,value) <- broadcastMap.value) yield { (x,key) } }} result.saveAsTextFile("file:///home/cjohnson/output")
补充说明
你当前使用广播变量关联小表的方案是合理的,只要rdd2的数据量可以单节点加载,该方案相比普通shuffle join可以大幅减少shuffle开销,完全适配大数据量的生产场景,无需调整广播逻辑。
内容的提问来源于stack exchange,提问作者weak_at_math
相关产品推荐
相关产品推荐

