【问题标题】:Spark and Python use custom file format/generator as input for RDDSpark 和 Python 使用自定义文件格式/生成器作为 RDD 的输入
【发布时间】:2014-10-05 21:30:51
【问题描述】:

我想问一下 Spark 中的输入可能性。我可以从http://spark.apache.org/docs/latest/programming-guide.html 看到,我可以使用sc.textFile() 将文本文件读取到 RDD,但我想在分发到 RDD 之前进行一些预处理,例如我的文件可能是 JSON 格式,例如. {id:123, text:"...", value:6} 我只想使用 JSON 的某些字段进行进一步处理。

我的想法是是否有可能以某种方式使用 Python 生成器作为 SparkContext 的输入?

或者,如果 Spark 中有一些更自然的方式如何处理自定义文件,而不是 Spark 的纯文本文件?

编辑:

似乎接受的答案应该有效,但它让我转向了我更实际的以下问题Spark and Python trying to parse wikipedia using gensim

【问题讨论】:

  • 您始终可以将 JSON 加载到 RDD 中,然后在 RDD 上进行处理以仅过滤您需要的数据。这样做的好处是这种“预处理”类型的工作可以在 Spark 集群中并行化。您能否举例说明您首先要进行哪种处理?
  • 通过预处理,我主要是指我只想从 JSON 或 XML 中选择例如字段 text1、text2,然后我可以做一些事情,比如用空格分割它并将其保存为文本文件。我没有看到任何自然的方式来解析 JSON RDD。现在我只能考虑将 JSON 或 XML 作为 sc.textFile() 文件进行处理,并且无论何时看到所需的密钥,然后使用以下字符串。你是这个意思吗?
  • 是的,我想我对你想要做的事情有感觉。如果我误解了,请评论我的回答。

标签: python hadoop apache-spark


【解决方案1】:

执行此操作的最快方法可能是按原样加载文本文件并进行处理以在生成的 RDD 上选择所需的字段。这使整个集群的工作并行化,并且比在单台机器上进行任何预处理更有效地扩展。

对于 JSON(甚至 XML),我认为您不需要自定义输入格式。由于 PySpark 在 Python 环境中执行,因此您可以使用 Python 中定期可用的函数来反序列化 JSON 并提取所需的字段。

例如:

import json

raw = sc.textFile("/path/to/file.json")
deserialized = raw.map(lambda x: json.loads(x))
desired_fields = deserialized.map(lambda x: x['key1'])

desired_fields 现在是原始 JSON 文件中key1 下所有值的 RDD。

您可以使用此模式来提取字段组合,用空格或其他方式将它们分割。

desired_fields = deserialized.map(lambda x: (x['key1'] + x['key2']).split(' '))

如果这变得太复杂,您可以将 lambda 替换为常规 Python 函数,该函数执行您想要的所有预处理,然后调用 deserialized.map(my_preprocessing_func)

【讨论】:

  • 我正在尝试实现和运行这个示例,所以一旦我让它工作,我会告诉你它是否对我有用。我现在很惊讶这实际上应该起作用,因为例如我的 JSON 可能大于所有机器上的 RAM 内存,所以我希望一次只加载它的一部分,因此 json.loads (x) 只会得到其中的一部分,并且无法正确解析它?
  • 实际上这应该可行,但它让我想到了stackoverflow.com/questions/26202978/…中的特定后续问题
  • @ziky90 “我的 JSON 可能比所有机器上的 RAM 内存都大” 这应该不是问题。如果源数据集中的每一行都有一个 JSON 对象,Spark 将能够很好地处理它,即使数据集很大。即使您尝试缓存数据集,Spark 也会尽可能多地缓存并从磁盘处理其余部分。
【解决方案2】:

是的,您可以使用 SparkContext.parallelize() 从 python 变量创建 RDD:

data = [1, 2, 3, 4, 5]
distData = sc.parallelize(data)
distData.count()   # 5

这个变量也可以是一个迭代器。

【讨论】:

  • 不幸的是,当生成器生成的列表大于 RAM 内存时,这将不起作用,因为 Spark 在内部将生成器转换为列表。您对这里的一些解决方法有什么想法吗?
  • 我担心会是这样。抱歉,除了在 Spark 中进行所有处理或尝试修补源代码之外,我不知道如何继续进行。
猜你喜欢
  • 2015-11-01
  • 2021-08-20
  • 2019-07-05
  • 1970-01-01
  • 1970-01-01
  • 2016-02-25
  • 1970-01-01
  • 2015-05-08
  • 2021-06-09
相关资源
最近更新 更多