【问题标题】:Computing custom function on numpy array results in UnpicklingError: NEWOBJ class argument has NULL tp_new在 numpy 数组上计算自定义函数会导致 UnpicklingError: NEWOBJ class argument has NULL tp_new
【发布时间】:2016-06-30 21:59:35
【问题描述】:

我遇到了一个非常奇怪的问题,在哪里使用辣.空间距离矩阵计算工作正常,但使用距离矩阵的自定义函数会导致 Spark 错误。

我的数据如下所示:

33.848366,-84.3733852,A,1234
33.848237299999994,-84.37318470000001,A,1234
33.8488057,-84.3731556,A,1234
33.847644200000005,-84.3727751,A,1234
33.84840429999999,-84.3732269,A,1234
33.849072899999996,-84.37342070000001,A,1234
33.8428191,-84.38306340000001,A,1234
33.842778499999994,-84.3830113,A,1234
33.8394582,-84.3770177,A,1234
33.847117299999994,-84.365351,A,1234

我的完全可重现的代码如下所示:

from pyspark import SparkContext

import pandas as pd
import numpy as np

from sklearn.cluster import DBSCAN
from math import radians, cos, sin, asin, sqrt
from scipy.spatial.distance import pdist, squareform

# This function taken from another StackOverflow post (modified radius only)
def distHaversine(pos1, pos2, r = 6378137):
    pos1 = pos1 * np.pi / 180
    pos2 = pos2 * np.pi / 180
    cos_lat1 = np.cos(pos1[..., 0])
    cos_lat2 = np.cos(pos2[..., 0])
    cos_lat_d = np.cos(pos1[..., 0] - pos2[..., 0])
    cos_lon_d = np.cos(pos1[..., 1] - pos2[..., 1])
    return r * np.arccos(cos_lat_d - cos_lat1 * cos_lat2 * (1 - cos_lon_d))

def myFunc(x):
    points = pd.DataFrame(list(x[1]))
    points.columns = ['lat', 'lon']
    ## PROBLEM LINE: UNCOMMENTING THIS LINE AND COMMENT BELOW TWO RESULTS IN THE ERROR ##
    # pointsDistMatrix = distHaversine(np.array(points)[:, None], np.array(points))
    pointsDistMatrix = pdist(points)
    pointsDistMatrix = squareform(pointsDistMatrix)
    db = DBSCAN(eps = 75, min_samples = 3, metric = 'precomputed',
                algorithm = 'kd_tree').fit(pointsDistMatrix)
    points['cluster'] = db.labels_
    return ((x[0], [tuple(x) for x in points.values]))

textFile = sc.textFile('df.csv')
processedGeoData = textFile \
                   .map(lambda x: x.split(',')) \
                   .map(lambda x: ((str(x[3]), str(x[2])),
                                        (float(x[0]), float(x[1])))) \
                   .groupByKey() \
                   .sortByKey(False) \
                   .map(myFunc)

processedGeoData.collect()

我得到的错误是:

Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.collectAndServe.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 50.0 failed 1 times, most recent failure: Lost task 0.0 in stage 50.0 (TID 70, localhost): org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/usr/local/Cellar/apache-spark/1.5.2/libexec/python/lib/pyspark.zip/pyspark/worker.py", line 98, in main
    command = pickleSer._read_with_length(infile)
  File "/usr/local/Cellar/apache-spark/1.5.2/libexec/python/lib/pyspark.zip/pyspark/serializers.py", line 164, in _read_with_length
    return self.loads(obj)
  File "/usr/local/Cellar/apache-spark/1.5.2/libexec/python/lib/pyspark.zip/pyspark/serializers.py", line 422, in loads
    return pickle.loads(obj)
UnpicklingError: NEWOBJ class argument has NULL tp_new

知道发生了什么吗?为什么自定义矩阵不起作用,但 scipy.spatial 起作用?

以下是我正在使用的不同软件包的版本:

Python 2.7.10
numpy==1.9.2
pandas==0.16.0rc1-22-g96aa9cb
scikit-learn==0.15.2
scipy==0.15.1
pyspark=1.5.2

【问题讨论】:

  • 当我使用scipy.spatial.distance.pdist(基本上是欧几里得距离)使用带有“预计算”选项的相同精确调用时,它工作正常。注释一行并取消注释其他两行会使其运行良好。如果 DBSCAN 调用有问题,不知道为什么会这样。是否存在一些包版本差异?
  • 顺便说一句-在我的情况下,我没有得到那个 ValueError 。改为获取此 UnpickingError。
  • 能否包含包和 Python 版本?
  • 已编辑问题并在末尾添加。谢谢!
  • 更新了所有软件包,现在正在使用 Anaconda。现在,我得到同样的 kd_tree 错误。我从DBSCAN 调用中删除了该选项并运行它。与pdist 一起工作正常,而它通过调用distHaversine 而不是pdist 生成TypeError: can't pickle ellipsis objects

标签: python numpy apache-spark scipy pyspark


【解决方案1】:

要使其正常工作,请创建一个模块 haversine.py:

import pandas as pd
import numpy as np

def distHaversine(pos1, pos2, r = 6378137):
    pos1 = pos1 * np.pi / 180
    pos2 = pos2 * np.pi / 180
    cos_lat1 = np.cos(pos1[..., 0])
    cos_lat2 = np.cos(pos2[..., 0])
    cos_lat_d = np.cos(pos1[..., 0] - pos2[..., 0])
    cos_lon_d = np.cos(pos1[..., 1] - pos2[..., 1])
    return r * np.arccos(cos_lat_d - cos_lat1 * cos_lat2 * (1 - cos_lon_d))

并分发它 (--py-files / sc.addPyFile)。下次导入distHaversine

>>> from haversine import distHaversine

然后你就解决了。

【讨论】:

  • 感谢您的协助。它确实帮助我缩小了问题的范围。请参阅下面的内容,了解无需拆分为单独文件的工作原理。
【解决方案2】:

好的,我不喜欢回答我自己的问题,但由于某些奇怪的原因,其他人遇到了同样的问题,我希望它会有所帮助。

我重写了上面的函数,不使用“省略号”运算符来索引数组,如下所示:

def distHaversine(pos1, pos2, r = 6378137):
    pos1 = pos1 * np.pi / 180
    pos2 = pos2 * np.pi / 180
    cos_lat1 = np.cos(pos1[:, :, 0])
    cos_lat2 = np.cos(pos2[:, 0])
    cos_lat_d = np.cos(pos1[:, :, 0] - pos2[:, 0])
    cos_lon_d = np.cos(pos1[:, :, 1] - pos2[:, 1])
    return r * np.arccos(cos_lat_d - cos_lat1 * cos_lat2 * (1 - cos_lon_d))

更改后一切正常。

似乎出于某种未知原因的省略号运算符正在干扰 Spark 序列化层?不知道幕后发生了什么,因为我什至没有将函数的返回值返回给 Spark。与标记为 pandas 数据框的集群完全不同的元组列表。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-07-05
    • 1970-01-01
    • 2017-03-27
    • 1970-01-01
    • 1970-01-01
    • 2012-10-12
    • 2019-09-30
    • 1970-01-01
    相关资源
    最近更新 更多