【问题标题】:pyspark: Save schemaRDD as json filepyspark:将 schemaRDD 保存为 json 文件
【发布时间】:2014-12-31 11:13:08
【问题描述】:

我正在寻找一种以 JSON 格式将数据从 Apache Spark 导出到各种其他工具的方法。我想一定有一个非常简单的方法来做到这一点。

示例:我有以下 JSON 文件 'jfile.json':

{"key":value_a1, "key2":value_b1},
{"key":value_a2, "key2":value_b2},
{...}

文件的每一行都是一个 JSON 对象。使用

可以轻松地将这类文件读入 PySpark
jsonRDD = jsonFile('jfile.json')

然后看起来像(通过调用 jsonRDD.collect()):

[Row(key=value_a1, key2=value_b1),Row(key=value_a2, key2=value_b2)]

现在我想将这些类型的文件保存回纯 JSON 文件。

我在 Spark 用户列表中找到了这个条目:

http://apache-spark-user-list.1001560.n3.nabble.com/Updating-exising-JSON-files-td12211.html

声称使用

RDD.saveAsTextFile(jsonRDD) 

这样做之后,文本文件看起来像

Row(key=value_a1, key2=value_b1)
Row(key=value_a2, key2=value_b2)

,即 jsonRDD 刚刚被明确写入文件。在阅读 Spark 用户列表条目后,我会期待一种“自动”转换回 JSON 格式。我的目标是有一个看起来像开头提到的“jfile.json”的文件。

我是否错过了一个非常明显的简单方法?

我阅读了http://spark.apache.org/docs/latest/programming-guide.html,在 google、用户列表和堆栈溢出中搜索了答案,但几乎所有答案都涉及读取 JSON 并将其解析为 Spark。我什至买了《Learning Spark》这本书,但那里的示例(第 71 页)只是导致与上面相同的输出文件。

有人可以帮我吗?我觉得我在这里只缺少一个小链接

提前干杯和感谢!

【问题讨论】:

    标签: python json apache-spark


    【解决方案1】:

    我看不出一个简单的方法来做到这一点。一种解决方案是将SchemaRDD 的每个元素转换为String,以RDD[String] 结尾,其中每个元素都针对该行格式化为JSON。因此,您需要编写自己的 JSON 序列化程序。那是容易的部分。它可能不是超级快,但应该可以并行工作,而且您已经知道如何将RDD 保存到文本文件中。

    关键的见解是,您可以通过调用schema 方法从SchemaRDD 中获得模式的表示。然后,map 传递给您的每个Row 都需要结合模式递归遍历。这实际上是平面 JSON 的串联列表遍历,但您可能还需要考虑嵌套 JSON。

    剩下的只是 Python 的小事,我不会说,但我确实有这个working in Scala,以防它帮助你。 Scala 代码变得密集的部分实际上并不依赖于深入的 Spark 知识,因此如果您能够理解基本递归并了解 Python,您应该能够使其工作。您的大部分工作是弄清楚如何在 Python API 中使用 pyspark.sql.Rowpyspark.sql.StructType

    请注意:我很确定我的代码在缺少值的情况下还不能工作——formatItem 方法需要处理空元素。

    编辑:Spark 1.2.0 中,toJSON 方法被引入SchemaRDD,这使得这个问题更加更简单 - - 查看@jegordon 的答案。

    【讨论】:

    【解决方案2】:

    您可以使用 toJson() 方法,它允许您将 SchemaRDD 转换为 JSON 文档的 MappedRDD。

    https://spark.apache.org/docs/latest/api/python/pyspark.sql.html?highlight=tojson#pyspark.sql.SchemaRDD.toJSON

    【讨论】:

    • 这很好用,而且是单线:res.toJSON().saveAsTextFile('/tmp/out/')
    【解决方案3】:

    我一直在从 SQL 控制台直接在 Spark SQL 中使用org.apache.spark.sql.json。这不是最有效的方法,它可能被认为是一种 hack,但它可以完成工作。

    CREATE TABLE jsonTable (
        key STRING,
        value STRING
    )
    USING org.apache.spark.sql.json
    OPTIONS (
        PATH "destination/path"
    );
    

    创建表后,从已注册的临时表或任何其他表中插入数据

    INSERT OVERWRITE TABLE jsonTable
    SELECT * FROM tempTable;
    

    注意:似乎这是在启动一个 hive map reduce 作业,在提供的路径下创建多个文件部分。预计执行缓慢

    注意:建表时提供​​的路径是hdfs,不是本地文件系统。

    注意:我没有尝试使用 SQLContext.sql 将其嵌入到脚本中,但它可能是可行的

    注意:从表 jsonTable 中选择可能会由于序列化而失败

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-07-09
      • 2019-04-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-30
      • 2018-07-31
      • 2021-11-26
      相关资源
      最近更新 更多