【问题标题】:Spark: How to "reduceByKey" when the keys are numpy arrays which are not hashable?Spark:当键是不可散列的numpy数组时,如何“reduceByKey”?
【发布时间】:2016-09-21 15:29:01
【问题描述】:

我有一个 (key,value) 元素的 RDD。键是 NumPy 数组。 NumPy 数组不可散列,当我尝试执行 reduceByKey 操作时,这会导致问题。

有没有办法用我的手动哈希函数提供 Spark 上下文?或者有没有其他方法可以解决这个问题(除了实际“离线”散列数组并将散列键传递给 Spark)?

这是一个例子:

import numpy as np
from pyspark import SparkContext

sc = SparkContext()

data = np.array([[1,2,3],[4,5,6],[1,2,3],[4,5,6]])
rd = sc.parallelize(data).map(lambda x: (x,np.sum(x))).reduceByKey(lambda x,y: x+y)
rd.collect()

错误是:

调用时出错 z:org.apache.spark.api.python.PythonRDD.collectAndServe。

...

TypeError: unhashable type: 'numpy.ndarray'

【问题讨论】:

    标签: python numpy pyspark rdd


    【解决方案1】:

    最简单的解决方案是将其转换为可散列的对象。例如:

    from operator import add
    
    reduced = sc.parallelize(data).map(
        lambda x: (tuple(x), x.sum())
    ).reduceByKey(add)
    

    如果需要,稍后再转换回来。

    有没有办法用我的手动哈希函数提供 Spark 上下文

    不是一个简单的。整个机制依赖于事实对象实现__hash__ 方法并且C 扩展不能被猴子修补。您可以尝试使用调度来覆盖pyspark.rdd.portable_hash,但我怀疑即使您考虑转换成本也是值得的。

    【讨论】:

      猜你喜欢
      • 2019-01-13
      • 2018-04-08
      • 1970-01-01
      • 2015-05-04
      • 1970-01-01
      • 1970-01-01
      • 2015-06-12
      • 1970-01-01
      • 2016-02-25
      相关资源
      最近更新 更多