PySpark如何将JSON字符串格式的字典列拆分为独立字段列
报错原因
你写的代码有两个核心错误:
- UDF返回类型定义不匹配:你声明的返回类型是
map<string, bigint>,但attributes字段解析后的值是字符串数组,正确的返回类型应为map<string, array<string>>,类型不匹配会触发序列化异常。 - 异常分支返回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
相关产品推荐
相关产品推荐

