You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 02:15:38