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

如何在PySpark中将JSON数组解析为多列并展开数据?

PySpark 实现JSON数组展开与字段提取

核心思路

要完成需求转换,需执行三个关键操作:

  • 解析job_details列的JSON数组为Spark可识别的结构化数组
  • 将数组拆分为多行(每行对应数组中的一个JSON对象)
  • 提取嵌套字段并替换原job_id,保留job_score和job_status字段

具体实现代码

1. 定义JSON Schema

根据job_details的结构定义匹配的Spark Schema:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, ArrayType

# 初始化SparkSession
spark = SparkSession.builder.appName("JsonArrayTransform").getOrCreate()

# 定义job_details的Schema
job_details_schema = ArrayType(
    StructType([
        StructField("jkl", StringType(), nullable=False),
        StructField("mno", StructType([
            StructField("uvw", StringType(), nullable=False)
        ]), nullable=False),
        StructField("stu", DoubleType(), nullable=False)
    ])
)

2. 模拟原始数据(可选,用于测试)

# 模拟原始表数据
data = [
    ("job_a", 234.5, "Active", """
        [
            {"jkl": "XYZDatafeed110m2211V1", "mno": {"uvw": "XYZ_TAG_30201.PV"}, "stu": 0.1234},
            {"jkl": "XYZDatafeed110m2211V1", "mno": {"uvw": "XYZ_TAG_30202.PV"}, "stu": 0.3623},
            {"jkl": "XYZDatafeed110m2211V1", "mno": {"uvw": "XYZ_TAG_30203.PV"}, "stu": 0.2345},
            {"jkl": "XYZDatafeed110m2211V1", "mno": {"uvw": "XYZ_TAG_30204.PV"}, "stu": 0.9345},
            {"jkl": "XYZDatafeed110m2211V1", "mno": {"uvw": "XYZ_TAG_30201.PV"}, "stu": 0.5345}
        ]
    """)
]

df = spark.createDataFrame(data, ["job_id", "job_score", "job_status", "job_details"])

3. 执行转换逻辑

from pyspark.sql.functions import from_json, explode, col

# 解析JSON字符串为结构化数组
parsed_df = df.withColumn("job_details_parsed", from_json(col("job_details"), job_details_schema))

# 将数组展开为多行
exploded_df = parsed_df.select("job_score", "job_status", explode(col("job_details_parsed")).alias("detail"))

# 提取目标字段并生成最终表
final_df = exploded_df.select(
    col("detail.jkl").alias("job_id"),
    col("job_score"),
    col("job_status"),
    col("detail.mno.uvw").alias("uvw"),
    col("detail.stu").alias("stu")
)

# 查看转换结果
final_df.show(truncate=False)

关键函数说明

  • from_json:将字符串格式的JSON数据解析为Spark结构化数据,需指定匹配的Schema确保解析准确性
  • explode:将数组列的每个元素拆分为单独行,实现"一行转多行"的展开效果
  • col:引用DataFrame列,配合.操作符直接提取嵌套层级的字段

内容的提问来源于stack exchange,提问作者MMV

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 07:19:59