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

PySpark嵌套字典转StringType格式异常问题求助

问题解答

这是正常行为,不是Bug

你遇到的格式变化是因为定义的columns字段为MapType(StringType, StringType),要求所有value必须是字符串类型。当嵌套字典/列表被存入这个Map时,PySpark会调用Python对象的__str__方法将其转为字符串,而不是标准JSON格式——这就导致了冒号变等号、引号消失的现象(比如{'a':1}转成{a=1}),这种格式不属于合法JSON,自然无法被from_json()解析。

修复方案

针对动态schema的场景,提供两种可行的修复思路:

思路1:直接存储完整的原始JSON字符串

如果你的API返回的是原始JSON字符串,不要提前解析为字典,直接将整个JSON字符串存入DataFrame的StringType字段,后续通过自动推断schema来解析:

import pyspark.sql.types as T
import json
from pyspark.sql.functions import from_json, schema_of_json

# 模拟API返回的原始JSON字符串(如果已经解析为字典,可通过json.dumps转回)
project_data = [
    (1, json.dumps({"issue_key":"abc", "col1":{'a':2342,'b':2}})),
    (2, json.dumps({"issue_key":"abc", "col2":[{'a':2342,'b':2}, {'c':'grha','d':True}]})),
    (3, json.dumps({"issue_key":"abc", "col1":{'a':5,'b':21}}))
]

# 定义schema,columns为StringType
schema = T.StructType([
    T.StructField("project_id", T.IntegerType()),
    T.StructField("columns", T.StringType())
])

df = spark.createDataFrame(project_data, schema)

# 自动推断JSON的schema
sample_json = df.select("columns").first()[0]
inferred_schema = schema_of_json(sample_json)

# 解析JSON为结构化数据
df_parsed = df.withColumn("columns_parsed", from_json("columns", inferred_schema))
df_parsed.display()

思路2:将嵌套结构转为标准JSON字符串后存入Map

如果需要保留MapType结构,可在存入DataFrame前,将嵌套的字典/列表通过json.dumps()转为标准JSON字符串,确保Map的value是合法JSON格式:

import pyspark.sql.types as T
import json
from pyspark.sql.functions import from_json, schema_of_json

# 原始数据
project_data = [
    (1,{"issue_key":"abc", "col1":{'a':2342,'b':2}}),
    (2,{"issue_key":"abc", "col2":[{'a':2342,'b':2}, {'c':'grha','d':True}]}),
    (3,{"issue_key":"abc", "col1":{'a':5,'b':21}})
]

# 处理函数:将嵌套结构转为JSON字符串
def convert_to_json_str(obj):
    if isinstance(obj, (dict, list)):
        return json.dumps(obj)
    return str(obj)

# 预处理数据,把每个Map的value转成JSON字符串
processed_data = [
    (proj_id, {k: convert_to_json_str(v) for k, v in cols.items()})
    for proj_id, cols in project_data
]

schema = T.StructType([
    T.StructField("project_id", T.IntegerType()),
    T.StructField("columns", T.MapType(T.StringType(), T.StringType()))
])

df = spark.createDataFrame(processed_data, schema)

# 示例:解析col1字段
sample_col1 = df.select("columns.col1").first()[0]
col1_schema = schema_of_json(sample_col1)
df_parsed = df.withColumn("col1_parsed", from_json(df["columns"]["col1"], col1_schema))
df_parsed.display()

关键说明

  • 动态schema场景下,schema_of_json()可以通过样本JSON字符串自动推断schema,无需硬编码结构。
  • 避免直接将非字符串类型的嵌套结构存入MapType(StringType, StringType),必须显式转为标准JSON字符串才能保证后续解析正常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 19:37:50