【问题标题】:PySpark reduceByKey? to add Key/TuplePySpark reduceByKey?添加键/元组
【发布时间】:2015-07-02 05:40:45
【问题描述】:

我有以下数据,我想做的是

[(13, 'D'), (14, 'T'), (32, '6'), (45, 'T'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'T'), (53, '2'), (54, '0'), (13, 'A'), (14, 'T'), (32, '6'), (45, 'A'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'X')]

是为每个键计数值的实例(1 个字符串字符)。所以我先做了一张地图:

.map(lambda x: (x[0], [x[1], 1]))

现在使它成为一个键/元组:

[(13, ['D', 1]), (14, ['T', 1]), (32, ['6', 1]), (45, ['T', 1]), (47, ['2', 1]), (48, ['0', 1]), (49, ['2', 1]), (50, ['0', 1]), (51, ['T', 1]), (53, ['2', 1]), (54, ['0', 1]), (13, ['A', 1]), (14, ['T', 1]), (32, ['6', 1]), (45, ['A', 1]), (47, ['2', 1]), (48, ['0', 1]), (49, ['2', 1]), (50, ['0', 1]), (51, ['X', 1])]

我只是无法在最后一部分弄清楚如何为每个键计算该字母的实例。例如,密钥 13 将有 1 个 D 和 1 个 A。而 14 将有 2 个 T,等等。

【问题讨论】:

  • 您要先groupByKey,然后对已分组的字符进行计数。

标签: python apache-spark pyspark


【解决方案1】:

我对 Scala 中的 Spark 更加熟悉,因此可能有比 Counter 更好的方法来计算 groupByKey 生成的迭代中的字符,但这里有一个选项:

from collections import Counter

rdd = sc.parallelize([(13, 'D'), (14, 'T'), (32, '6'), (45, 'T'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'T'), (53, '2'), (54, '0'), (13, 'A'), (14, 'T'), (32, '6'), (45, 'A'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'X')]) 
rdd.groupByKey().mapValues(Counter).collect()

[(48, Counter({'0': 2})),
 (32, Counter({'6': 2})),
 (49, Counter({'2': 2})),
 (50, Counter({'0': 2})),
 (51, Counter({'X': 1, 'T': 1})),
 (53, Counter({'2': 1})),
 (13, Counter({'A': 1, 'D': 1})),
 (45, Counter({'A': 1, 'T': 1})),
 (14, Counter({'T': 2})),
 (54, Counter({'0': 1})),
 (47, Counter({'2': 2}))]

【讨论】:

  • 哦,你已经用过Counter了!不幸的是,应该避免使用groupByKey,因为它聚合了 master 上的所有数据。而且 2 次操作而不是 1 次也是不够的。但 1 票赞成紧凑!
  • @ipoteka 有趣的是,我不知道groupByKey 的效率低下你有没有详细说明的好参考?
  • 不错的链接,非常漂亮的数字。完全有道理。我怀疑有时groupByKey 仍然会“足够快”,但很高兴知道。
  • @Nikita 它不会汇总“master”上的所有数据。但它以非聚合的形式被改组给执行者。这就是 reduceByKey 的关键区别,它在对数据进行洗牌之前执行一个聚合步骤,因此(通常)通过网络发送的数据要少得多。
【解决方案2】:

代替:

.map(lambda x: (x[0], [x[1], 1]))

我们可以这样做:

.map(lambda x: ((x[0], x[1]), 1))

在最后一步,我们可以使用 reduceByKeyadd。请注意,add 来自 operator 包。

把它放在一起:

from operator import add
rdd = sc.parallelize([(13, 'D'), (14, 'T'), (32, '6'), (45, 'T'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'T'), (53, '2'), (54, '0'), (13, 'A'), (14, 'T'), (32, '6'), (45, 'A'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'X')]) 
rdd.map(lambda x: ((x[0], x[1]), 1)).reduceByKey(add).collect()

【讨论】:

    【解决方案3】:

    如果我没听错的话,你可以一次操作combineByKey

    from collections import Counter
    x = sc.parallelize([(13, 'D'), (14, 'T'), (32, '6'), (45, 'T'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'T'), (53, '2'), (54, '0'), (13, 'A'), (14, 'T'), (32, '6'), (45, 'A'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'X')]) 
    result = x.combineByKey(lambda value:  {value: 1}, 
    ...                     lambda x, value:  value.get(x,0) + 1,
    ...                     lambda x, y: dict(Counter(x) + Counter(y)))
    result.collect()
    [(32, {'6': 2}), (48, {'0': 2}), (49, {'2': 2}), (53, {'2': 1}), (13, {'A': 1, 'D': 1}), (45, {'A': 1, 'T': 1}), (50, {'0': 2}), (54, {'0': 1}), (14, {'T': 2}), (51, {'X': 1, 'T': 1}), (47, {'2': 2})]
    

    【讨论】:

    • 看起来这个解决方案 13 有 ('A', 2) 而不是 [('A', 1), ('D', 1)]
    • 嗯,我假设 13 只对应于“A”,我会改变我的答案。谢谢!
    • OP 希望对与每个键关联的每个字符进行计数
    • @ohruunuruus 我编辑了它,但我不确定解决方案是否足够“pythonic”。
    • 我得到的一件事是:AttributeError: 'str' object has no attribute 'get'
    【解决方案4】:

    我尝试使用函数和mapValues() 转换

    def f(Counter): return Counter
    
    from collections import Counter
    
    rdd=sc.parallelize([(13, 'D'), (14, 'T'), (32, '6'), (45, 'T'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'T'), (53, '2'), (54, '0'), (13, 'A'), (14, 'T'), (32, '6'), (45, 'A'), (47, '2'), (48, '0'), (49, '2'), (50, '0'), (51, 'X')])
    rdd.groupByKey().mapValues(Counter).collect()
    

    【讨论】:

      猜你喜欢
      • 2015-10-17
      • 2016-12-27
      • 2015-06-25
      • 2021-10-26
      • 1970-01-01
      • 1970-01-01
      • 2015-12-09
      • 2018-08-07
      • 2017-05-17
      相关资源
      最近更新 更多