代码之家  ›  专栏  ›  技术社区  ›  Venkatesh

在pyspark中将json文件展平为单个行

  •  0
  • Venkatesh  · 技术社区  · 3 年前

    我收到了一个json文件作为api的输入,这里是示例json。

    json_data = 
    {
      "field1": "value1",
      "field2": "value2",
      "message_records": [
        {
          "field3": "value3",
          "field4": "value4"
        },
        {
          "field5": "value5",
          "field6": "value6"
        }
      ],
      "messages": [
        {
          "field7": "value3",
          "field8": "value4"
        },
        {
          "field9": "value5",
          "field10": "value6"
        },
        {
          "field11": "value5",
          "field12": "value6"
        }
      ]
    }
    

    如何使用python将json数据扁平化为单独的行,并将数据加载到dataframe中。这里,具有嵌套数组的messages、messagerecords需要加载到单独的记录中。 将json文件转换为pyspark数据帧

    这里,字段1和字段2对于message_records和消息是常见的。我需要将message_record数据写入单独的文件,并将消息数据写入单独文件

    0 回复  |  直到 3 年前
        1
  •  1
  •   JayashankarGS    3 年前

    您可以使用以下代码在单独的行中创建,并将数据写入的单独文件 消息记录 消息

    record = {}
    record["field1"] = json_data["field1"]
    record["field2"] = json_data["field2"]
    message_records_df =spark.createDataFrame([record])
    messages_df = spark.createDataFrame([record])
    

    使用创建两个数据帧 field1 field2 。由于这两者都与相同 消息记录 消息

    from pyspark.sql.types import LongType
    from pyspark.sql import Row
    
    def zipindexdf(df):
        schema_new = df.schema.add("index", LongType(), False)
        return df.rdd.zipWithIndex().map(lambda l: list(l[0]) + [l[1]]).toDF(schema_new)
        
    message_records_df_index = zipindexdf(message_records_df)
    message_records_df_index.show()
    messages_df_index = zipindexdf(messages_df)
    messages_df_index.show()
    

    我在这里添加 index 列使用 zipWithIndex 所以它可以在索引上联接。

    enter image description here

    接下来,创建数据帧并通过循环遍历中的每个项合并到最终数据帧 消息记录 如下所示。

    for i in json_data['message_records']:
        df = zipindexdf(spark.createDataFrame([i]))
        message_records_df_index = message_records_df_index.join(df, "index", "inner")
        message_records_df_index.show()
    

    enter image description here

    和我一样 消息 如下所示

    for i in json_data['messages']:
        df = zipindexdf(spark.createDataFrame([i]))
        messages_df_index = messages_df_index.join(df, "index", "inner")
        messages_df_index.show()
    

    enter image description here

    最后,将这些数据写入csv文件。

    message_records_df_index.write.option("header","true").csv('/message_records_df/')
    messages_df_index.write.option("header","true").csv('/messages_df/')
    
    spark.read.option("header","true").csv('/message_records_df/').show()
    spark.read.option("header","true").csv('/messages_df/').show()
    

    enter image description here

        2
  •  1
  •   Anupam Chand    3 年前

    你可以在中找到答案 this SO answer

    你只需要改变你对段落的称呼。我用过

    df_flat1 = flatten_test(multiline_df.select(multiline_df.field1,multiline_df.field2,multiline_df.message_records))
    df_flat2 = flatten_test(multiline_df.select(multiline_df.field1,multiline_df.field2, multiline_df.messages))
    df_flat1.printSchema()
    df_flat1.show(5)
    
    df_flat2.printSchema()
    df_flat2.show(5)
    

    并且得到

    root
     |-- field1: string (nullable = true)
     |-- field2: string (nullable = true)
     |-- message_records_field3: string (nullable = true)
     |-- message_records_field4: string (nullable = true)
     |-- message_records_field5: string (nullable = true)
     |-- message_records_field6: string (nullable = true)
    
    +------+------+----------------------+----------------------+----------------------+----------------------+
    |field1|field2|message_records_field3|message_records_field4|message_records_field5|message_records_field6|
    +------+------+----------------------+----------------------+----------------------+----------------------+
    |value1|value2|                value3|                value4|                  null|                  null|
    |value1|value2|                  null|                  null|                value5|                value6|
    +------+------+----------------------+----------------------+----------------------+----------------------+
    
    root
     |-- field1: string (nullable = true)
     |-- field2: string (nullable = true)
     |-- messages_field10: string (nullable = true)
     |-- messages_field11: string (nullable = true)
     |-- messages_field12: string (nullable = true)
     |-- messages_field7: string (nullable = true)
     |-- messages_field8: string (nullable = true)
     |-- messages_field9: string (nullable = true)
    
    +------+------+----------------+----------------+----------------+---------------+---------------+---------------+
    |field1|field2|messages_field10|messages_field11|messages_field12|messages_field7|messages_field8|messages_field9|
    +------+------+----------------+----------------+----------------+---------------+---------------+---------------+
    |value1|value2|            null|            null|            null|         value3|         value4|           null|
    |value1|value2|          value6|            null|            null|           null|           null|         value5|
    |value1|value2|            null|          value5|          value6|           null|           null|           null|
    +------+------+----------------+----------------+----------------+---------------+---------------+---------------+
    
    推荐文章