【问题标题】:creating pair rdd from two rdds based repetition count of first rdd in pyspark?从pyspark中第一个rdd的两个基于重复计数的rdds创建对rdd?
【发布时间】:2019-05-06 13:15:57
【问题描述】:

我已经创建了 2 个如下所示的 Rdd

rd2=sc.parallelize([['A', 'B','D'], ['B', 'C'], ['A', 'B'],['B']])
rd3=sc.parallelize([['A', 'B'],['B', 'C'],['B','D']])
rd2.collect()
[['A', 'B','D'], ['B', 'C'], ['A', 'B'],['B']]
rd3.collect()
[['A', 'B'], ['B', 'C'],['B','D']]

现在我想将两个 rdd 在 rd2 中的重复计数中的公共元素计算为新 rd4 中的值,即

['A', 'B'] 在两个 rdd 中都很常见,但在 rd2 中重复计数为 2。

我预期的 rd4 是:

[(['A','B'],2),(['B','C'],1),(['B','D'],1)]

【问题讨论】:

  • 如果你可以使用数据框而不是 RDD,这将是一个简单的连接,然后是聚合计数
  • 你是对的,但根据我的要求,我不应该使用 DF

标签: python apache-spark dataframe pyspark rdd


【解决方案1】:

您可以在 rdd 上使用 countByKey 检查重复计数,它将返回一个 defaultdict

但你说你希望你的结果为rdd,所以你可以使用reduceByKey函数。

我会创建和你一样的rdd

rd2=sc.parallelize([['A', 'B'], ['B', 'C'], ['A', 'B'],['B']])

rd2.map(lambda x: (tuple(x),1)).reduceByKey(lambda x,y: x+y).collect()
[(('B',), 1), (('A', 'B'), 2), (('B', 'C'), 1)]

现在您的输出 rdd 为 (tuple,count) 结构,您可以通过 map 函数将其更改为列表。

rd2.map(lambda x: (tuple(x),1)).reduceByKey(lambda x,y: x+y).map(lambda x: (list(x[0]),x[1])).collect()
[(['B'], 1), (['A', 'B'], 2), (['B', 'C'], 1)] 

希望这能解决您的问题。

【讨论】:

  • 根据我的问题,我想将 rdd 的 rd3,rd2.repetition 计数中的共同元素计算为新 rd4 rdd.but 根据您的代码,您不考虑 rd3.above万一它会失败。请帮助我。提前致谢。
  • 那么您的预期输出是什么? @Sai
  • 根据我的问题,我们应该考虑两个 rdds,但我们不考虑 rd3。仅基于 rd3 元素,我们应该计算 rd2 元素。请看一下我已经更新了我的问题以便于理解
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-08-19
  • 1970-01-01
  • 2020-09-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-12-10
相关资源
最近更新 更多