PySpark:在RDD中使用Object

Dan*_*iel 7 python apache-spark pyspark

我目前正在学习Python,并希望将其应用于/使用Spark.我有这个非常简单(无用)的脚本:

import sys
from pyspark import SparkContext

class MyClass:
    def __init__(self, value):
        self.v = str(value)

    def addValue(self, value):
        self.v += str(value)

    def getValue(self):
        return self.v

if __name__ == "__main__":
    if len(sys.argv) != 1:
        print("Usage CC")
        exit(-1)

    data = [1, 2, 3, 4, 5, 2, 5, 3, 2, 3, 7, 3, 4, 1, 4]
    sc = SparkContext(appName="WordCount")
    d = sc.parallelize(data)
    inClass = d.map(lambda input: (input, MyClass(input)))
    reduzed = inClass.reduceByKey(lambda a, b: a.addValue(b.getValue))
    print(reduzed.collect())
Run Code Online (Sandbox Code Playgroud)

用它执行时

spark-submit CustomClass.py

..以下错误是thorwn(输出缩短):

Caused by: org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/usr/local/spark/python/lib/pyspark.zip/pyspark/worker.py", line 111, in main
    process()
  File "/usr/local/spark/python/lib/pyspark.zip/pyspark/worker.py", line 106, in process
    serializer.dump_stream(func(split_index, iterator), outfile)
  File "/usr/local/spark/python/lib/pyspark.zip/pyspark/serializers.py", line 133, in dump_stream
    for obj in iterator:
  File "/usr/local/spark/python/lib/pyspark.zip/pyspark/rdd.py", line 1728, in add_shuffle_key
  File "/usr/local/spark/python/lib/pyspark.zip/pyspark/serializers.py", line 415, in dumps
    return pickle.dumps(obj, protocol)
PicklingError: Can't pickle __main__.MyClass: attribute lookup __main__.MyClass failed
at org.apache.spark.api.python.PythonRunner$$anon$1.read(PythonRDD.scala:166)...
Run Code Online (Sandbox Code Playgroud)

声明给我

PicklingError: Can't pickle __main__.MyClass: attribute lookup __main__.MyClass failed

似乎很重要.这意味着类实例无法序列化,对吧?你知道如何解决这个问题吗?

感谢致敬

Kev*_*n S 16

有很多问题:

  • 如果你放入MyClass一个单独的文件,它可以被腌制.这是许多Python使用pickle的常见问题.这很容易通过移动MyClass和使用来解决from myclass import MyClass.通常dill可以修复这些问题(如import dill as pickle),但这对我不起作用.
  • 一旦解决了这个问题,你的reduce就不起作用了,因为调用addValuereturn None(不返回),而不是一个实例MyClass.你需要改变addValue回来self.
  • 最后,lambda需要打电话getValue,所以应该有a.addValue(b.getValue())

一起: myclass.py

class MyClass:
    def __init__(self, value):
        self.v = str(value)

    def addValue(self, value):
        self.v += str(value)
        return self

    def getValue(self):
        return self.v
Run Code Online (Sandbox Code Playgroud)

main.py

import sys
from pyspark import SparkContext
from myclass import MyClass

if __name__ == "__main__":
    if len(sys.argv) != 1:
        print("Usage CC")
        exit(-1)

    data = [1, 2, 3, 4, 5, 2, 5, 3, 2, 3, 7, 3, 4, 1, 4]
    sc = SparkContext(appName="WordCount")
    d = sc.parallelize(data)
    inClass = d.map(lambda input: (input, MyClass(input)))
    reduzed = inClass.reduceByKey(lambda a, b: a.addValue(b.getValue()))
    print(reduzed.collect())
Run Code Online (Sandbox Code Playgroud)