【问题标题】:Spark inner join two RDD and new value should be a sum of the old valuesSpark内部连接两个RDD,新值应该是旧值的总和
【发布时间】:2017-04-01 14:11:53
【问题描述】:

我有两个 RDD,我想要一个内连接,那么新值应该是 rdd1rdd2 中的值的总和

一个例子:

RDD1 = [(key1,(1,2)), (key2,(2,3)),(key3,(3,4))]
RDD2 = [(key1,(2,3), (key3,(4,5))]

RDD_inner_join= [(key1,(3,5), (key3,(7,9))]

我尝试过使用rdd1.join(rdd2),但我想当我们加入两个rdd时,我们无法对旧值进行任何操作。

有人有办法解决我的问题吗?

提前致谢。

【问题讨论】:

  • join 是个好方法,加入后你就可以访问 RDD 元素了。能否通过 join 分享您的代码,以便我们帮助您找到其中的错误?
  • 谢谢,我同意加入是一个好方法,但是加入后我们有一个像这样的元组(key,(value1,value2)),但是我需要一个像(key,value)加入后的元组+value2)

标签: python apache-spark inner-join rdd


【解决方案1】:

加入两个 RDD 后,执行如下映射转换:

joinedRDDs.map(lambda r: (r[0], r[1][0][0] + r[1][1][0], r[1][0][1] + r[1][1][1]))

你会得到预期的结果。

【讨论】:

  • 谢谢 Mariusz,如果我们在加入后做另一个工作 mapreduce,那将是个好主意。
  • 实际上,在 spark 中没有“mapreduce job”之类的东西。您的应用程序将执行两个转换(混洗连接和非混洗映射),然后执行一些操作(例如将结果保存到存储中,显示它等)。没关系,基本上就是你在 spark 中编写应用程序的方式。
猜你喜欢
  • 1970-01-01
  • 2018-08-09
  • 1970-01-01
  • 1970-01-01
  • 2016-08-26
  • 2017-09-04
  • 2018-04-08
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多