PySpark如何将列中嵌套JSON字符串解析拆分为指定多列
PySpark 解析JSON字符串列提取嵌套字段实现方法
你的columnName列为STRING类型,存储标准嵌套JSON结构,外层为data节点,data下包含直接存储的业务字段,以及嵌套的address地址对象,共需提取8个字段,以下是两种可直接落地的实现方案:
方案1:使用get_json_object快速提取(适合快速开发、临时查询场景)
该方法无需提前定义JSON结构,直接按JSON路径提取字段,路径规则为$.层级名.字段名,嵌套字段顺着层级写路径即可。
from pyspark.sql import functions as F # 假设原始DataFrame变量名为df parsed_df = df.select( "*", F.get_json_object("columnName", "$.data.address.city").alias("city"), F.get_json_object("columnName", "$.data.address.state").alias("state"), F.get_json_object("columnName", "$.data.address.street1").alias("street1"), F.get_json_object("columnName", "$.data.address.zipCode").alias("zipCode"), F.get_json_object("columnName", "$.data.cardAlias").alias("cardAlias"), F.get_json_object("columnName", "$.data.storeName").alias("storeName"), F.get_json_object("columnName", "$.data.cardNumber").alias("cardNumber"), F.get_json_object("columnName", "$.data.storeNumber").alias("storeNumber") )
该方案优缺点:
- 优点:写法简单,不需要提前对齐JSON结构,JSON存在小幅字段变动不会直接抛错
- 缺点:每提取一个字段就要对JSON字符串做一次解析,数据量大、提取字段多时性能较差;所有提取结果默认是字符串类型,需要数值、日期等其他类型必须额外手动转换
方案2:使用from_json定义Schema一次性解析(生产环境推荐,性能更优)
该方法提前定义和JSON结构匹配的Schema,只做一次JSON解析就可以提取所有字段,性能更高,定义Schema时可以直接指定字段类型,省去后续转换步骤,方便维护。
from pyspark.sql import functions as F from pyspark.sql.types import ( StructType, StructField, StringType, IntegerType ) # 定义和JSON层级完全匹配的Schema,字段名大小写必须和原始JSON一致 json_schema = StructType([ StructField("data", StructType([ # 嵌套的address地址节点 StructField("address", StructType([ StructField("city", StringType()), StructField("state", StringType()), StructField("street1", StringType()), StructField("zipCode", StringType()) ])), # data节点下的直接业务字段,需要数值类型可以直接替换类型定义,比如storeNumber用IntegerType StructField("cardAlias", StringType()), StructField("storeName", StringType()), StructField("cardNumber", StringType()), StructField("storeNumber", StringType()) ])) ]) parsed_df = df.withColumn("parsed_json", F.from_json("columnName", json_schema)) \ .select( "*", "parsed_json.data.address.city", "parsed_json.data.address.state", "parsed_json.data.address.street1", "parsed_json.data.address.zipCode", "parsed_json.data.cardAlias", "parsed_json.data.storeName", "parsed_json.data.cardNumber", "parsed_json.data.storeNumber" ) \ .drop("parsed_json") # 解析完成后删除临时结构化列
小提示:如果原始列存在不合法的JSON脏数据,from_json解析后对应行的结构化字段会返回null,不会导致整个任务失败,可以后续加过滤规则单独处理脏数据。
结果验证
解析完成后执行parsed_df.show(truncate=False)即可查看提取结果,如果出现字段为null的情况,优先检查JSON路径、Schema定义的字段名和原始JSON的拼写、大小写是否完全匹配。
内容的提问来源于stack exchange,提问作者Kishor
相关产品推荐
相关产品推荐

