如何使用PySpark提取嵌套JSON中的指定字段?
用PySpark提取JSON中的指定字段
我有一个带有如下Schema的JSON文件:
root |-- context: struct (nullable = true) | |-- application: struct (nullable = true) | | |-- version: string (nullable = true) | |-- custom: struct (nullable = true) | | |-- dimensions: array (nullable = true) | | | |-- element: struct (containsNull = true) | | | | |-- Activity ID: string (nullable = true) | | | | |-- Activity Type: string (nullable = true) | | | | |-- Bot ID: string (nullable = true) | | | | |-- Channel ID: string (nullable = true) | | | | |-- Conversation ID: string (nullable = true) | | | | |-- Correlation ID: string (nullable = true) | | | | |-- From ID: string (nullable = true) | | | | |-- Recipient ID: string (nullable = true) | | | | |-- StatusCode: string (nullable = true) | | | | |-- Timestamp: string (nullable = true) | |-- data: struct (nullable = true) | | |-- eventTime: string (nullable = true) | | |-- isSynthetic: boolean (nullable = true) | | |-- samplingRate: double (nullable = true) | |-- device: struct (nullable = true) | | |-- roleInstance: string (nullable = true) | | |-- roleName: string (nullable = true) | | |-- type: string (nullable = true) | |-- location: struct (nullable = true) | | |-- city: string (nullable = true) | | |-- clientip: string (nullable = true) | | |-- continent: string (nullable = true) | | |-- country: string (nullable = true) | | |-- province: string (nullable = true) | |-- operation: struct (nullable = true) | | |-- id: string (nullable = true) | | |-- parentId: string (nullable = true) | |-- session: struct (nullable = true) | | |-- isFirst: boolean (nullable = true) |-- event: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- count: long (nullable = true) | | |-- name: string (nullable = true) |-- internal: struct (nullable = true) | |-- data: struct (nullable = true) | | |-- documentVersion: string (nullable = true) | | |-- id: string (nullable = true)
JSON文件示例:
{ "event": [ { "name": "Activity", "count": 1 } ], "internal": { "data": { "id": "79baca55-d168-11ea-b166-6fc861e9e21c", "documentVersion": "1.61" } }, "context": { "application": { "version": "Wed 07/22/2020 5:37:05.58 \r\nUTC (fv-az461) [Build 148886] [Repo Intercom] [Branch prod] [Commit XXX] \r\n[XX 1.6.20-140775] [XXX 1.3.27-144047] \r\n" }, "data": { "eventTime": "2020-07-29T06:55:15.6294636Z", "isSynthetic": false, "samplingRate": 100 }, "cloud": {}, "device": { "type": "PC", "roleName": "bc-directline-southindia", "roleInstance": "RD0003FF905CCA", "screenResolution": {} }, "session": { "isFirst": false }, "operation": { "id": "XXX", "parentId": "|XXXX.c4cd9570_" }, "location": { "clientip": "0.0.0.0", "continent": "XX", "country": "XXX", "province": "XXX", "city": "XXX" }, "custom": { "dimensions": [ { "Timestamp": "XXX" }, { "StatusCode": "200" }, { "Activity ID": "JoH4veTvChCCnzchOD1Lg-f|0000001" }, { "From ID": "XXX" }, { "Correlation ID": "|54734cb21ba7f143a72ddd03fc865669.c4cd9570_" }, { "Channel ID": "directline" }, { "Recipient ID": "XXXX" }, { "Bot ID": "XXXX" }, { "Activity Type": "message" }, { "Conversation ID": "XXX" } ] } } }
需要提取Activity ID、Activity Type、Bot ID、Channel ID、Conversation ID、Correlation ID、From ID、Recipient ID、StatusCode、Timestamp到DataFrame,实现步骤如下:
实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, map_keys, map_values, first, monotonically_increasing_id # 初始化SparkSession spark = SparkSession.builder.appName("ExtractTargetFields").getOrCreate() # 读取JSON文件 df = spark.read.json("path/to/your/json/files") # 给每个原始行添加唯一ID,确保后续透视后每行对应原始一个JSON对象 df_with_row_id = df.withColumn("row_id", monotonically_increasing_id()) # 展开context.custom.dimensions数组,保留行ID exploded_df = df_with_row_id.select( "row_id", explode(col("context.custom.dimensions")).alias("dimension") ) # 将单个键的结构体转换为字段名和字段值的行 key_value_df = exploded_df.select( "row_id", map_keys(col("dimension"))[0].alias("field_name"), map_values(col("dimension"))[0].alias("field_value") ) # 按行ID分组,透视字段名为列,提取每个字段对应的值 final_df = key_value_df.groupBy("row_id").pivot("field_name").agg(first("field_value")) # 移除临时的row_id列 final_df = final_df.drop("row_id") # 查看结果 final_df.show()
代码说明
- 添加行ID:因为每个原始JSON的
dimensions是数组,每个元素对应一个字段,添加唯一ID可以保证透视后每行对应原始一个JSON对象,不会出现字段值跨行聚合的问题。 - 展开数组:用
explode函数将数组拆分为多行,每个数组元素对应一行记录。 - 提取键值对:利用
map_keys和map_values从单个键的结构体中取出字段名和对应的值。 - 透视转列:通过
pivot将字段名转为列,再用first聚合每个行ID下的字段值(每个行ID下每个字段只会出现一次)。 - 清理临时列:移除用于分组的
row_id,得到仅包含目标字段的DataFrame。
内容的提问来源于stack exchange,提问作者Alexander M
相关产品推荐
相关产品推荐

