【问题标题】:Pyspark: create a schema from JSON filePyspark:从 JSON 文件创建模式
【发布时间】:2021-12-11 14:37:15
【问题描述】:

我正在处理来自非常长的嵌套 JSON 文件的数据。问题是,这些文件的结构并不总是与其中一些遗漏其他列的相同。我想从包含所有列的空 JSON 文件创建自定义架构。如果我稍后将 JSON 文件读入此预定义模式,则不存在的列将填充空值(至少是计划)。到目前为止我做了什么:

  1. 将测试 JSON(不包含预期的所有列)加载到数据帧中
  2. 将其架构写入 JSON 文件
  3. 在文本编辑器中打开此 JSON 文件并手动添加缺少的列

接下来我想做的是通过将 JSON 文件读入我的代码来创建一个新模式,但我在使用 synthax 时遇到了困难。我可以直接从文件本身读取架构吗?我试过了

schemaFromJson = StructType.fromJson(json.loads('filepath/spark-schema.json'))

但它给了我 TypeError: init() missing 2 required positional arguments: 'doc' 和 'pos'

知道我当前的代码有什么问题吗? 非常感谢

编辑: 我遇到了这个链接 sparkbyexamples.com/pyspark/pyspark-structtype-and-structfield 。第 7 章几乎描述了我遇到的问题。我只是不明白如何解析我手动增强为 schemaFromJson = StructType.fromJson(json.loads(schema.json)) 的 json 文件。

当我这样做时:

jsonDF = spark.read.json(filesToLoad)
schema = jsonDF.schema.json()
schemaNew = StructType.fromJson(json.loads(schema))
jsonDF2 = spark.read.schema(schemaNew).json(filesToLoad)

代码运行通过,但显然没有用,因为 jsonDF 和 jsonDF2 确实具有相同的内容/架构。我想要实现的是在“schema”中添加一些列,然后这些列将反映在“schemaNew”中。

【问题讨论】:

    标签: pyspark apache-spark-sql jsonschema


    【解决方案1】:

    您为什么不定义一个空的 DF,其中包含 JSON 文件可以具有的所有列?然后将 JSON 加载到其中。这是一个想法:

    对于 Spark 3.1.0:

    from pyspark.sql.types import *
    
    schema = StructType([
        StructField("fruit",StringType(),True),
        StructField("size",StringType(),True),
        StructField("color",StringType(),True)
    ])
    df = spark.createDataFrame([], schema)
    
    json_file_1 = {"fruit": "Apple","size": "Large"}
    json_df_1 = spark.read.json(sc.parallelize([json_file_1]))
    
    df = df.unionByName(json_df_1, allowMissingColumns=True)
    
    json_file_2 = {"fruit": "Banana","size": "Small","color": "Yellow"}
    
    df = df.unionByName(json_file_2, allowMissingColumns=True)
    
    display(df)
    

    【讨论】:

    • 好吧,我试过了,但这是一场噩梦,因为我必须定义 300 多列,其中大部分是嵌套的。我遇到了这个链接sparkbyexamples.com/pyspark/pyspark-structtype-and-structfield。第 7 章几乎描述了我遇到的问题。我只是不明白如何解析我手动增强为 schemaFromJson = StructType.fromJson(json.loads(schema.json)) 的 json 文件
    【解决方案2】:

    我想我明白了。 Schemapath 包含已增强的架构:

    schemapath = '/path/spark-schema.json'
    with open(schemapath) as f:
       d = json.load(f)
       schemaNew = StructType.fromJson(d)
       jsonDf2 = spark.read.schema(schmaNew).json(filesToLoad)
       jsonDF2.printSchema()
    

    【讨论】:

      猜你喜欢
      • 2016-08-09
      • 2017-11-19
      • 1970-01-01
      • 1970-01-01
      • 2019-05-11
      • 2022-01-17
      • 2015-11-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多