【问题标题】:PySpark - UnpicklingError: NEWOBJ class argument has NULL tp_newPySpark - UnpicklingError: NEWOBJ 类参数有 NULL tp_new
【发布时间】:2017-06-06 17:12:03
【问题描述】:

当我执行波纹管时,我收到 Unpickling Error

rdd = sc.parallelize([('HOMICIDE', {'2017': 1}), 
('DECEPTIVE PRACTICE', {'2015': 2, '2017': 2, '2016': 8}), 
('ROBBERY', {'2016': 2})])

rdd.flatMapValues(dict.items).collect()

错误如下,在字典值上使用flatMapValues有什么问题

  File "/usr/hdp/2.3.4.0-3485/spark/python/lib/pyspark.zip/pyspark/worker.py", line 98, in main
    command = pickleSer._read_with_length(infile)
  File "/usr/hdp/2.3.4.0-3485/spark/python/lib/pyspark.zip/pyspark/serializers.py", line 164, in _read_with_length
    return self.loads(obj)
  File "/usr/hdp/2.3.4.0-3485/spark/python/lib/pyspark.zip/pyspark/serializers.py", line 422, in loads
    return pickle.loads(obj)
UnpicklingError: NEWOBJ class argument has NULL tp_new
) [duplicate 3]
17/06/06 17:01:14 INFO TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool 

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:
    rdd = sc.parallelize([('HOMICIDE', {'2017': 1}), 
                          ('DECEPTIVE PRACTICE', {'2015': 2, '2017': 2, '2016': 8}), 
                          ('ROBBERY', {'2016': 2})])
    
    rdd.flatMapValues(lambda data: data.items()).collect()
    
    [('HOMICIDE', ('2017', 1)),
     ('DECEPTIVE PRACTICE', ('2015', 2)),
     ('DECEPTIVE PRACTICE', ('2017', 2)),
     ('DECEPTIVE PRACTICE', ('2016', 8)),
     ('ROBBERY', ('2016', 2))]
    

    dict.items 是方法描述符。您必须提供一个函数来通知 flatmap 如何解压缩这些值。我通过将 labmda 函数传递给 flatMap 函数来做到这一点。

    【讨论】:

      猜你喜欢
      • 2013-07-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-04-06
      • 2020-12-20
      • 2019-04-27
      • 2017-11-14
      • 2021-07-14
      相关资源
      最近更新 更多