如何在PySpark及本地环境中将JSON列展开为多列?
问题描述
需求:将数据中的JSON对象列(Column B)展开为多列。
原始数据:
| Column A | Column B |
|---|---|
| id1 | [{a:1,b:'letter1'}] |
| id2 | [{a:1,b:'letter2',c:3,d:4}] |
期望转换结果:
| Column A | a | b | c | d |
|---|---|---|---|---|
| id1 | 1 | letter1 | ||
| id2 | 1 | letter2 | 3 | 4 |
遇到的问题:
- 本地Pandas环境:提取键值对后转DataFrame时触发
ValueError: All arrays must be of the same length,原因是部分JSON对象缺少部分键。 - PySpark环境:Pandas转Spark DataFrame时出现类型不兼容问题;转为字符串后,执行substring操作提示
Column is not iterable;尝试转换JSON字符串为JSON对象时出现ValueError: 'json' is not in list ; AttributeError: json错误。且因存在大量不同JSON列,无法硬编码Schema。
本地Pandas解决方案
直接用pd.json_normalize处理,自动兼容缺失键的情况,步骤如下:
import pandas as pd import ast # 原始数据(如果Column B是字符串格式,先转成Python字典) df = pd.DataFrame({ 'Column A': ['id1', 'id2'], 'Column B': ["[{a:1,b:'letter1'}]", "[{a:1,b:'letter2',c:3,d:4}]"] }) # 字符串转字典列表 df['Column B'] = df['Column B'].apply(lambda x: ast.literal_eval(x)) # 提取列表中的第一个字典并归一化 normalized_df = pd.json_normalize(df['Column B'].apply(lambda x: x[0])) # 合并原数据与归一化结果 result_df = pd.concat([df['Column A'], normalized_df], axis=1) print(result_df)
解释:json_normalize会自动为缺失的键填充NaN,不会触发长度不匹配的错误。如果Column B本来就是字典列表格式,可跳过ast.literal_eval的转换步骤。
PySpark解决方案
针对无固定Schema的场景,通过动态解析JSON生成列,步骤如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_replace, from_json, map_keys, array_distinct, col, explode from pyspark.sql.types import StringType, MapType spark = SparkSession.builder.appName("JSONExpand").getOrCreate() # 原始数据(Column B为字符串类型) data = [ ("id1", "[{a:1,b:'letter1'}]"), ("id2", "[{a:1,b:'letter2',c:3,d:4}]") ] df = spark.createDataFrame(data, ["Column A", "Column B"]) # 第一步:将非标准JSON转为标准格式(给键添加双引号,移除列表括号) df = df.withColumn("json_str", regexp_replace("Column B", r"(\w+):", r'"\1":')) df = df.withColumn("json_str", regexp_replace("json_str", r"^\[|\]$", "")) # 第二步:转为MapType并提取所有唯一键 df = df.withColumn("json_map", from_json(col("json_str"), MapType(StringType(), StringType()))) all_keys = df.select(explode(map_keys(col("json_map")))).distinct().rdd.flatMap(lambda x: x).collect() # 第三步:动态生成所有列 for key in all_keys: df = df.withColumn(key, col("json_map").getItem(key)) # 保留目标列 result_df = df.select("Column A", *all_keys) result_df.show()
解释:
- 先将非标准JSON转为
from_json可识别的标准格式; - 通过
MapType灵活解析键值对,再提取所有可能的键; - 动态生成列,缺失的键自动填充
null,无需硬编码Schema。
若Column B本身是数组类型的字典,可简化为:
# 假设Column B为ArrayType(MapType)格式 df = df.withColumn("json_dict", explode(col("Column B"))) result_df = df.groupBy("Column A").pivot(map_keys(col("json_dict"))).agg(first(col("json_dict").getItem(map_keys(col("json_dict")))))
内容的提问来源于stack exchange,提问作者Jack9406
相关产品推荐
相关产品推荐

