【问题标题】:reduceByKey not working as expected when keys are of type bitarray当键是位数组类型时,reduceByKey 无法按预期工作
【发布时间】:2017-03-29 17:13:02
【问题描述】:

以下是我在 pyspark shell 中尝试的代码。

from bitarray import bitarray
a = bitarray('0') * 5
b = bitarray('1') * 5
c = [a.copy() for x in range(3)]
d = [b.copy() for x in range(5)]
e = c + d
rdd = sc.parallelize(e).map(lambda x : (x, 1)).reduceByKey(lambda x, y : x + y).collect()
print(rdd)

预期:

[(bitarray('11111'), 5), (bitarray('00000'), 3)

实际输出:

[(bitarray('11111'), 1), (bitarray('00000'), 1), (bitarray('11111'), 1), (bitarray('11111'), 1), (bitarray('00000'), 1), (bitarray('11111'), 1), (bitarray('00000'), 1), (bitarray('11111'), 1)]

为什么 spark 引擎不能区分不同的位数组值?

【问题讨论】:

  • a[0] 和 a[1] 的 hash 和 equals 是什么?我相信通过 id(a[0]) 显示,但我不是真正的 python 人
  • 不要认为这是可能的,看这里stackoverflow.com/a/29725381/4964651
  • 谢谢@mtoto!这验证了问题。哈希必须相等。

标签: python apache-spark pyspark


【解决方案1】:

规避该问题的一种方法是使用tobytes 将所有bitarray 转换为字节。另外,@mtoto 提供的链接中讨论了问题的原因,这是 bitarray 不保持哈希不变。

from bitarray import bitarray

def back2Bit(U):
        res = bitarray()
        res.frombytes(U)
        return res

a = bitarray('0') * 5
b = bitarray('1') * 5
c = [a.copy() for x in range(3)]
d = [b.copy() for x in range(5)]
e = c + d
rdd = sc.parallelize([x.tobytes() for x in e]).map(lambda x : (x, 1)).reduceByKey(lambda x, y : x + y).collect()
rdd = [(back2Bit(x[0]), x[1]) for x in rdd]

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-07-19
    • 2016-01-08
    • 2018-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多