【问题标题】:Spark DataFrame serialized as invalid jsonSpark DataFrame 序列化为无效 json
【发布时间】:2018-07-08 06:37:25
【问题描述】:

TL;DR:当我将 Spark DataFrame 转储为 json 时,我总是会得到类似的结果

{"key1": "v11", "key2": "v21"}
{"key1": "v12", "key2": "v22"}
{"key1": "v13", "key2": "v23"}

这是无效的 json。我可以手动编辑转储文件以获取可以解析的内容:

[
  {"key1": "v11", "key2": "v21"},
  {"key1": "v12", "key2": "v22"},
  {"key1": "v13", "key2": "v23"}
]

但我很确定我错过了一些可以让我避免手动编辑的东西。我只是现在不知道。

更多详情:

我有一个org.apache.spark.sql.DataFrame,我尝试使用以下代码将其转储到 json:

myDataFrame.write.json("file.json")

我也尝试过:

myDataFrame.toJSON.saveAsTextFile("file.json")

在这两种情况下,它最终都会正确转储每一行,但它缺少行之间的分隔逗号以及方括号。 因此,当我随后尝试解析此文件时,我使用的解析器会侮辱我,然后失败。

如果我能了解如何转储有效的 json,我将不胜感激。 (阅读DataFrameWriter 的文档并没有给我任何有趣的提示。)

【问题讨论】:

    标签: json apache-spark apache-spark-sql spark-dataframe


    【解决方案1】:

    这是预期的输出。 Spark 使用JSON Lines-like 格式有很多原因:

    • 可以并行解析和加载。
    • 无需在内存中加载完整文件即可完成解析。
    • 可以并行编写。
    • 无需在内存中存储完整分区即可写入。
    • 即使文件为空也是有效的输入。
    • 最后,Spark 中的 Row 是映射到 JSON 对象而不是数组的结构。
    • ...

    您可以通过几种方式创建所需的输出,但它总是会与上述任何一种冲突。

    例如,您可以为每个分区编写一个 JSON 文档:

    import org.apache.spark.sql.functions._
    
    df
      .groupBy(spark_partition_id)
      .agg(collect_list(struct(df.columns map col: _*)).alias("data"))
      .select($"data")
      .write
      .json(output_path)
    

    您可以在此前面加上 repartition(1) 以获得单个输出文件,但这不是您想要做的事情,除非数据非常小。

    1.6 替代方案将是glom

    import org.apache.spark.sql.Row
    import org.apache.spark.sql.types._
    
    val newSchema = StructType(Seq(StructField("data", ArrayType(df.schema))))
    
    sqlContext.createDataFrame(
      df.rdd.glom.flatMap(a => if(a.isEmpty) Seq() else Seq(Row(a))), 
      newSchema
    )
    

    【讨论】:

    猜你喜欢
    • 2019-12-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多