如何使用PySpark展开JSON列?附数据示例与预期结果
解析PySpark DataFrame中的嵌套JSON字符串
原始数据结构
你的PySpark DataFrame包含两个字段:value(存储嵌套JSON格式的字符串)和timestamp,数据样例如下:
+-----------------------------------------------------------------------------------------------+-----------------------+ |value |timestamp | +-----------------------------------------------------------------------------------------------+-----------------------+ |{"after":{"id":1001,"first_name":"Sally","last_name":"Thomas","email":"sally.thomas@acme.com"}}|2023-01-03 11:02:11.975| |{"after":{"id":1002,"first_name":"George","last_name":"Bailey","email":"gbailey@foobar.com"}} |2023-01-03 11:02:11.976| |{"after":{"id":1003,"first_name":"Edward","last_name":"Walker","email":"ed@walker.com"}} |2023-01-03 11:02:11.976| |{"after":{"id":1004,"first_name":"Anne","last_name":"Kretchmar","email":"annek@noanswer.org"}} |2023-01-03 11:02:11.976| +-----------------------------------------------------------------------------------------------+-----------------------+
Schema信息:
root |-- value: string (nullable = true) |-- timestamp: timestamp (nullable = true)
实现步骤及代码
要得到扁平化的结果,需要解析value字段中的JSON字符串,提取嵌套的after结构并展开为单独列,具体代码如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json from pyspark.sql.types import StructType, StructField, IntegerType, StringType # 初始化SparkSession(如果未初始化) spark = SparkSession.builder.appName("ParseNestedJSON").getOrCreate() # 定义JSON字符串对应的Schema(推荐显式定义,避免自动推断的不确定性) json_schema = StructType([ StructField("after", StructType([ StructField("id", IntegerType(), nullable=True), StructField("first_name", StringType(), nullable=True), StructField("last_name", StringType(), nullable=True), StructField("email", StringType(), nullable=True) ]), nullable=True) ]) # 1. 解析value字段为结构体类型 df_parsed = df.withColumn("value_struct", from_json(df.value, json_schema)) # 2. 提取after字段并展开为单独列 result_df = df_parsed.select( "value_struct.after.id", "value_struct.after.first_name", "value_struct.after.last_name", "value_struct.after.email" # 若需保留timestamp字段,添加 ", timestamp" 即可 ) # 更简洁写法:先提取after列,再展开所有子字段 # result_df = df_parsed.select("value_struct.after.*") # 查看结果 result_df.show()
代码说明
- 显式定义Schema:提前定义JSON对应的结构体Schema,避免自动推断可能带来的解析误差,适配生产环境的稳定性要求。
- 解析JSON字符串:用
from_json函数将value列的字符串转为PySpark结构体类型,生成新列value_struct。 - 提取展开字段:通过
.操作符访问嵌套字段,直接选择目标列;或用.*一键展开after下的所有子字段,简化代码。
执行后得到预期的扁平化DataFrame:
+-----+----------+---------+-----------------------+ | id|first_name|last_name| email| +-----+----------+---------+-----------------------+ |1001| Sally| Thomas|sally.thomas@acme.com | |1002| George| Bailey| gbailey@foobar.com | |1003| Edward| Walker| ed@walker.com | |1004| Anne| Kretchmar| annek@noanswer.org | +-----+----------+---------+-----------------------+
内容的提问来源于stack exchange,提问作者losforword
相关产品推荐
相关产品推荐

