【发布时间】:2017-05-11 10:04:44
【问题描述】:
我有一个 rdd 说 sample_rdd 类型为 RDD[(String, String, Int))] 有 3 列 id、item、count。样本数据:
id1|item1|1
id1|item2|3
id1|item3|4
id2|item1|3
id2|item4|2
我想将每个 id 加入到 lookup_rdd 这个:
item1|0
item2|0
item3|0
item4|0
item5|0
输出应该为我提供以下 id1、带有查找表的外连接:
item1|1
item2|3
item3|4
item4|0
item5|0
同样对于 id2 我应该得到:
item1|3
item2|0
item3|0
item4|2
item5|0
每个 id 的最终输出应该具有 id 的所有计数:
id1,1,3,4,0,0
id2,3,0,0,2,0
重要提示:此输出应始终按照查找中的顺序进行排序
这是我尝试过的:
val line = rdd_sample.map { case (id, item, count) => (id, (item,count)) }.map(row=>(row._1,row._2)).groupByKey()
get(line).map(l=>(l._1,l._2)).mapValues(item_count=>lookup_rdd.leftOuterJoin(item_count))
def get (line: RDD[(String, Iterable[(String, Int)])]) = { for{ (id, item_cnt) <- line i = item_cnt.map(tuple => (tuple._1,tuple._2)) } yield (id,i)
【问题讨论】:
-
val line = rdd_sample.map { case (id, item, count) => (id, (item,count)) }.map(row=>(row._1,row._2)).groupByKey() -
get(line).map(l=>(l._1,l._2)).mapValues(item_count=>lookup_rdd.leftOuterJoin(item_count))函数:def get (line: RDD[(String, Iterable[(String, Int)])]) = { for{ (id, item_cnt) <- line i = item_cnt.map(tuple => (tuple._1,tuple._2)) } yield (id,i) } -
您可以将其编辑到问题中。
-
@NanditaDwivedi 你试过解决方案了吗?
标签: scala apache-spark spark-dataframe