【问题标题】:spark: access rdd inside another rddspark:在另一个 rdd 中访问 rdd
【发布时间】: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


【解决方案1】:

坏消息是你不能在另一个 RDD 中使用。

好消息是,对于您的用例,假设 6000 个条目相当小,有一个理想的解决方案:收集驱动程序上的 RDD,将其广播回集群的每个节点并在另一个节点中使用它和以前一样的 RDD。

val sc: SparkContext = ???
val headers = sc.broadcast(lookup_rdd.zipWithIndex.collect().toMap)
val indexing = data_rdd.map { case (_, item_cnt_list ) =>
  item_cnt_list.map { case (k, v) => (headers.value(k), v) }
}
indexing.collect().foreach(println)

【讨论】:

  • 感谢您的回答。有类似的类型情况,但另外..必须更新 map 函数中的查找表。对于下一个元素,我必须在更新的查找表上进行查找。我知道我们不能用广播做到这一点。你能建议如何处理这个问题。即使是资源的链接也会有所帮助。提前致谢。
  • 我相信你有更好的改变为你的特定案例创建一个问题,分享相关代码。没有它很难说。
  • 添加了一个单独的问题:你能看看吗。 :stackoverflow.com/questions/49125735/…
  • 我明白整个答案。我有一个与此链接中相同的问题:您能看看吗? (stackoverflow.com/questions/60522369/…)
猜你喜欢
  • 1970-01-01
  • 2016-12-23
  • 1970-01-01
  • 1970-01-01
  • 2017-08-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多