【发布时间】:2017-05-15 10:49:19
【问题描述】:
我有一个大小为 6000 的查找 rdd,lookup_rdd: RDD[String]
a1
a2
a3
a4
a5 .....
还有另一个 rdd,data_rdd: RDD[(String, Iterable[(String, Int)])]: (id,(item,count)) 具有唯一的 id,
(id1,List((a1,2), (a3,4)))
(id2,List((a2,1), (a4,2), (a1,1)))
(id3,List((a5,1)))
lookup_rdd 中的 FOREACH 元素我想检查每个 id 是否有该元素,如果有,我输入计数,如果没有,我输入 0,并存储在文件中。
实现这一目标的有效方法是什么。可以散列吗?例如。我想要的输出是:
id1,2,0,4,0,0
id2,1,1,0,2,0
id3,0,0,0,0,1
我试过这个:
val headers = lookup_rdd.zipWithIndex().persist()
val indexing = data_rdd.map{line =>
val id = line._1
val item_cnt_list = line._2
val arr = Array.fill[Byte](6000)(0)
item_cnt_list.map(c=>(headers.lookup(c._1),c._2))
}
indexing.collect().foreach(println)
我得到了例外:
org.apache.spark.SparkException: RDD transformations and actions can only be invoked by the driver, not inside of other transformations
【问题讨论】:
-
6000 个数据集是一个非常小的数据集。考虑在驱动上采集然后广播
标签: scala apache-spark apache-spark-sql spark-dataframe