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

PySpark如何将JSON字符串格式的字典列拆分为独立字段列

报错原因

你写的代码有两个核心错误:

  1. UDF返回类型定义不匹配:你声明的返回类型是map<string, bigint>,但attributes字段解析后的值是字符串数组,正确的返回类型应为map<string, array<string>>,类型不匹配会触发序列化异常。
  2. 异常分支返回None:你写的UDF遇到JSON解析错误时直接pass(返回None),explode无法处理None值,会直接抛出运行时错误。

原生PySpark实现方案

不需要转pandas,以下是高效的原生实现代码:

步骤1:定义正确的JSON解析UDF

from pyspark.sql import functions as F
from pyspark.sql.types import MapType, ArrayType, StringType
import json

@F.udf(returnType=MapType(StringType(), ArrayType(StringType())))
def parse_attr_json(s):
    if not s:
        return {}
    try:
        return json.loads(s)
    except json.JSONDecodeError:
        return {}

步骤2:解析attributes字段并生成结果表

两种方案可选,按需使用:

方案1:先收集全量键再生成列(适合大数据量,性能更高)
# 解析JSON为Map类型列
df_parsed = df.withColumn("attr_map", parse_attr_json("attributes"))
# 收集所有存在的属性键
all_attr_keys = df_parsed.select(F.explode("attr_map")).select("key").distinct().rdd.map(lambda row: row.key).collect()
# 逐个提取键对应的值作为新列
result_df = df_parsed.select(
    F.col("id"),
    *[F.col("attr_map").getItem(key).alias(key) for key in all_attr_keys]
)
方案2:用pivot实现(写法更简洁,适合中小数据量)
# 解析JSON为Map类型列
df_parsed = df.withColumn("attr_map", parse_attr_json("attributes"))
# 炸开Map后分组转宽表
result_df = df_parsed.select("id", F.explode("attr_map")) \
    .groupBy("id") \
    .pivot("key") \
    .agg(F.first("value"))

运行result_df.show(truncate=False)即可得到你需要的输出结果。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 02:24:01