【问题标题】:Pickling a Spark RDD and reading it into Python腌制 Spark RDD 并将其读入 Python
【发布时间】:2016-02-21 20:47:41
【问题描述】:

我正在尝试通过酸洗来序列化 Spark RDD,并将酸洗后的文件直接读入 Python。

a = sc.parallelize(['1','2','3','4','5'])
a.saveAsPickleFile('test_pkl')

然后我将 test_pkl 文件复制到我的本地。如何将它们直接读入 Python?当我尝试普通的泡菜包时,当我尝试读取“test_pkl”的第一个泡菜部分时它失败了:

pickle.load(open('part-00000','rb'))

Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/usr/lib64/python2.6/pickle.py", line 1370, in load
    return Unpickler(file).load()
  File "/usr/lib64/python2.6/pickle.py", line 858, in load
    dispatch[key](self)
  File "/usr/lib64/python2.6/pickle.py", line 970, in load_string
    raise ValueError, "insecure string pickle"
ValueError: insecure string pickle

我假设 spark 使用的酸洗方法与 python 酸洗方法不同(如果我错了,请纠正我)。有什么方法可以让我从 Spark 中提取数据并将这个腌制对象直接从文件中读取到 python 中?

【问题讨论】:

  • 问题是它不是一个泡菜文件,而是一个包含泡菜对象的SequenceFile,我不知道有任何积极开发的 Python 中的 SequenceFiles 解析器。

标签: python apache-spark pickle pyspark


【解决方案1】:

更好的方法可能是腌制每个分区中的数据,对其进行编码,然后将其写入文本文件:

import cPickle
import base64

def partition_to_encoded_pickle_object(partition):
    p = [i for i in partition] # convert the RDD partition to a list
    p = cPickle.dumps(p, protocol=2) # pickle the list
    return [base64.b64encode(p)] # base64 encode the list, and return it in an iterable

my_rdd.mapPartitions(partition_to_encoded_pickle_object).saveAsTextFile("your/hdfs/path/")

将文件下载到本地目录后,可以使用以下代码段读取:

# you first need to download the file, this step is not shown
# afterwards, you can use 
path = "your/local/path/to/downloaded/files/"
data = []
for part in os.listdir(path):
    if part[0] != "_": # this prevents system generated files from getting read - e.g. "_SUCCESS"
        data += cPickle.loads(base64.b64decode((open(part,'rb').read())))

【讨论】:

  • 这里唯一的问题是加载部分需要将所有数据加载到data的内存中,这可能并不总是可行的。
  • @Tgsmith61591 正确 - 如果您在单台机器上读取数据,您通常无法读取集群中的所有数据。要解决此问题,您可能希望仅从文件中过滤/缩小/提取所需的数据,例如data += some_filter_fx(cPickle.loads(base64.b64decode((open(part,'rb').read()))))
【解决方案2】:

问题是格式不是泡菜文件。它是腌制objects 的SequenceFile。 sequence file 可以在 Hadoop 和 Spark 环境中打开,但不打算在 python 中使用并使用基于 JVM 的序列化来序列化,在这种情况下是字符串列表。

【讨论】:

    【解决方案3】:

    可以使用sparkpickle 项目。就这么简单

    with open("/path/to/file", "rb") as f:
        print(sparkpickle.load(f))
    

    【讨论】:

      猜你喜欢
      • 2018-05-14
      • 2017-08-21
      • 2017-04-22
      • 2018-03-03
      • 2023-03-18
      • 1970-01-01
      • 1970-01-01
      • 2017-04-02
      • 1970-01-01
      相关资源
      最近更新 更多