我的代碼如下所示:
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)
val result = rdd1.map{case(x,y) => {
for((key,value) <- broadcastMap .value) {
(x,key)
}
}}
result.saveAsTextFile("file:///home/cjohnson/output")
寫入檔案的預期輸出應該是:
(1,2)
(1,3)
(1,7)
(2,2)
(2,3)
(2,7)
(3,2)
(3,3)
(3,7)
但我將此輸出寫入檔案:
()
()
()
我怎樣才能解決這個問題?
PS 這只是我提供的一些小樣本資料,以證明我的問題。實際資料要大得多。
uj5u.com熱心網友回復:
- 那個內部
for回傳 Unit 即()因為你忘記添加yield:
val a = for((key, _) <- Map(1 -> "")) yield { (key) }
- 您需要
flatMap而不是map在每個 rdd 鍵和廣播映射鍵之間制作該產品。
關于你的問題,這是我將如何處理它:
rdd1
.keys
.flatMap { rddKey =>
broadcastMap
.value
.keys
.map(broadcastKey => (rddKey, broadcastKey))
}
后期編輯:
它可以寫成笛卡爾
rdd1
.keys
.cartesian(rdd2.keys)
轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/372774.html
上一篇:如何使用AkkaHTTP解組json回應以洗掉不必要的欄位
下一篇:將包含值的列作為串列轉換為陣列
