SPARK-5063 与尝试嵌套 RDD 操作时更好的错误消息相关,这是不支持的。
这是一个可用性问题,而不是功能问题。根本原因是 RDD 操作的嵌套,解决方法是打破它。
在这里,我们正在尝试 dRDD 和 mRDD 的连接。如果mRDD 的大小很大,建议使用rdd.join,否则,如果mRDD 很小,即适合每个执行器的内存,我们可以收集它,广播它并做一个'map-side ' 加入。
加入
一个简单的连接应该是这样的:
val rdd = sc.parallelize(Seq(Array("one","two","three"), Array("four", "five", "six")))
val map = sc.parallelize(Seq("one" -> 1, "two" -> 2, "three" -> 3, "four" -> 4, "five" -> 5, "six"->6))
val flat = rdd.flatMap(_.toSeq).keyBy(x=>x)
val res = flat.join(map).map{case (k,v) => v}
如果我们想使用广播,我们首先需要在本地收集解析表的值,以便将其传递给所有执行者。 注意要广播的 RDD必须适合驱动程序和每个执行程序的内存。
Map-side JOIN 与广播变量
val rdd = sc.parallelize(Seq(Array("one","two","three"), Array("four", "five", "six")))
val map = sc.parallelize(Seq("one" -> 1, "two" -> 2, "three" -> 3, "four" -> 4, "five" -> 5, "six"->6)))
val bcTable = sc.broadcast(map.collectAsMap)
val res2 = rdd.flatMap{arr => arr.map(elem => (elem, bcTable.value(elem)))}