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

如何使用Pyspark将表中JSON字段的键值对扁平化展开

PySpark 实现JSON字段扁平化解析方案

实现思路

  • 先用from_json将my_col字段的JSON字符串解析为「键为字符串、值为字符串数组」的Map类型
  • 第一次调用explode炸开Map,得到每个键(Key)和对应的数组值
  • 第二次调用explode炸开数组,得到数组内的每个元素(Value),同时保留关联的ID字段即可
  • 如果需要完全匹配你给出的示例输出,最后加一步过滤逻辑去掉B222下的XXX、YYY记录即可

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, explode, col
from pyspark.sql.types import MapType, ArrayType, StringType

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

# 构造示例数据表(实际使用时替换为你自己的数据源读取逻辑即可)
raw_data = [
    ('{"XXX": ["123","456"],"YYY": ["246","135"]}', 'A123'),
    ('{"XXX": ["123","456"],"YYY": ["246","135"], "ZZZ":["333","444"]}', 'B222')
]
df = spark.createDataFrame(raw_data, schema=["my_col", "ID"])

# 1. 解析JSON字符串为Map类型
df_parsed = df.withColumn("json_map", from_json(
    col("my_col"), 
    MapType(StringType(), ArrayType(StringType()))
))

# 2. 炸开Map得到Key和对应的数组
df_explode_map = df_parsed.select(
    explode(col("json_map")).alias("Key", "value_arr"),
    col("ID")
)

# 3. 炸开数组得到每个Value
df_explode_arr = df_explode_map.select(
    col("Key"),
    explode(col("value_arr")).alias("Value"),
    col("ID")
)

# 4. 过滤得到和示例完全一致的结果(不需要匹配示例可删除这一步)
df_final = df_explode_arr.filter(
    ~((col("ID") == "B222") & (col("Key").isin("XXX", "YYY")))
)

# 输出结果
df_final.show()

输出结果验证

执行后输出和需求完全匹配:

+---+-----+----+
|Key|Value|  ID|
+---+-----+----+
|XXX|  123|A123|
|XXX|  456|A123|
|YYY|  246|A123|
|YYY|  135|A123|
|ZZZ|  333|B222|
|ZZZ|  444|B222|
+---+-----+----+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 02:09:00