【问题标题】:How do I read a Large JSON Array File in PySpark如何在 PySpark 中读取大型 JSON 数组文件
【发布时间】:2018-07-21 10:14:29
【问题描述】:

问题

我最近在尝试读取大型 UTF-8 JSON 数组文件并切换到 HDInsight PySpark(v2.x,而不是 3)来处理文件时遇到了 Azure Data Lake Analytics 的挑战。该文件约为 110G,包含约 150m JSON 对象。

HDInsight PySpark 似乎不支持输入 JSON 文件格式的数组,所以我被卡住了。另外,我有“许多”这样的文件,每个文件都包含不同的架构,每个包含数百列,因此此时为这些文件创建架构不是一种选择。

问题

如何在 HDInsight 上的 PySpark 2 中使用开箱即用的功能来使这些文件能够以 JSON 格式读取?

谢谢,

J

我尝试过的事情

我使用了本页底部的方法: from Databricks 提供了以下代码 sn-p:

import json

df = sc.wholeTextFiles('/tmp/*.json').flatMap(lambda x: json.loads(x[1])).toDF()
display(df)

我尝试了上述方法,但不了解“wholeTextFiles”的工作原理,当然遇到了 OutOfMemory 错误,导致我的 executors 很快死亡。

我尝试加载到 RDD 和其他打开方法,但 PySpark 似乎只支持 JSONLines JSON 文件格式,并且由于 ADLA 对这种文件格式的要求,我有 JSON 对象数组。

我尝试以文本文件的形式读入,剥离 Array 字符,在 JSON 对象边界上拆分并像上面那样转换为 JSON,但这一直给出有关无法转换 unicode 和/或 str (ings) 的错误。

我找到了通过上述方法的方法,并将其转换为包含一列的数据框,其中包含 JSON 对象的字符串行。但是,我没有找到一种仅将数据框行中的 JSON 字符串输出到输出文件的方法。总是出来的

{'dfColumnName':'{...json_string_as_value}'}

我还尝试了一个 map 函数,它接受上述行,解析为 JSON,提取值(我想要的 JSON),然后将值解析为 JSON。这似乎可行,但是当我尝试保存时,RDD 是 PipelineRDD 类型并且没有 saveAsTextFile() 方法。然后我尝试了 toJSON 方法,但一直收到关于“找不到有效的 JSON 对象”的错误,我承认我不明白,当然还有其他转换错误。

【问题讨论】:

    标签: json azure pyspark rdd azure-hdinsight


    【解决方案1】:

    我终于找到了前进的方向。我了解到我可以直接从 RDD 中读取 json,包括 PipelineRDD。我找到了一种方法来删除 unicode 字节顺序标头,包装数组方括号,基于幸运分隔符拆分 JSON 对象,并拥有一个分布式数据集以进行更有效的处理。输出数据帧现在具有以 JSON 元素命名的列、推断架构并动态适应其他文件格式。

    这是代码 - 希望对您有所帮助!:

    #...Spark considers arrays of Json objects to be an invalid format
    #    and unicode files are prefixed with a byteorder marker
    #
    thanksMoiraRDD = sc.textFile( '/a/valid/file/path', partitions ).map(
        lambda x: x.encode('utf-8','ignore').strip(u",\r\n[]\ufeff") 
    )
    
    df = sqlContext.read.json(thanksMoiraRDD)
    

    【讨论】:

      猜你喜欢
      • 2023-04-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-04-05
      • 2014-05-19
      • 1970-01-01
      • 2020-01-08
      • 2021-12-17
      相关资源
      最近更新 更多