【发布时间】: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