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
相关产品推荐
相关产品推荐

