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

如何在PySpark及本地环境中将JSON列展开为多列?

问题描述

需求:将数据中的JSON对象列(Column B)展开为多列。
原始数据:

Column AColumn B
id1[{a:1,b:'letter1'}]
id2[{a:1,b:'letter2',c:3,d:4}]

期望转换结果:

Column Aabcd
id11letter1
id21letter234

遇到的问题:

  • 本地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()

解释:

  1. 先将非标准JSON转为from_json可识别的标准格式;
  2. 通过MapType灵活解析键值对,再提取所有可能的键;
  3. 动态生成列,缺失的键自动填充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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 21:50:14