PySpark创建DataFrame时如何强制保留JSON字符串格式?
解决PySpark中JSON字段被篡改(冒号变等号)的问题
核心原因
当你的Python字典中key10是Python字典对象而非字符串类型的JSON时,即使在Schema中指定StringType,Spark仍会将其序列化成key=value的格式(这是Spark对非字符串类型的默认toString行为),而非保留原始JSON结构。
最优解决方案:从源头处理数据
确保key10在传入Spark前已经是字符串类型的JSON,而非Python字典。步骤如下:
- 使用Python的
json.dumps()将key10的字典结构转为标准JSON字符串 - 再用处理后的数据创建DataFrame
示例代码:
import json from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType # 原始数据示例 data = [ { "key1": "value1", "key10": {"a": "b", "c": "data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mNk+M9QDwADhgGAWjR9awAAAABJRU5ErkJggg=="} } ] # 预处理:将key10转为JSON字符串 processed_data = [ {**item, "key10": json.dumps(item["key10"])} for item in data ] # 定义Schema schema = StructType([ StructField("key1", StringType(), nullable=True), StructField("key10", StringType(), nullable=True) ]) # 创建Spark会话并生成DataFrame spark = SparkSession.builder.appName("FixJSONField").getOrCreate() df = spark.createDataFrame(processed_data, schema) # 验证结果(不会出现冒号变等号的情况) df.show(truncate=False)
应急修复:处理已被篡改的DataFrame
如果已经生成了被篡改的DataFrame,可尝试用正则表达式精准替换键名后的等号(注意:此方法有局限性,仅适用于键名是字母数字的场景):
from pyspark.sql.functions import regexp_replace, concat, lit # 假设df是已被篡改的DataFrame,key10内容类似"a=b,c=data:image/png;base64,..." df_fixed = df.withColumn("key10_fixed", regexp_replace("key10", r"(\w+)=", r'$1:"')) \ .withColumn("key10_fixed", regexp_replace("key10_fixed", r",(\w+)=", r',$1:"')) \ .withColumn("key10_fixed", regexp_replace("key10_fixed", r"=(?!//|/)", r'":')) \ .withColumn("key10_fixed", concat(lit("{"), "key10_fixed", lit("}"))) # 查看修复结果 df_fixed.select("key10_fixed").show(truncate=False)
注意:此正则方法无法覆盖所有边缘情况(比如键名含特殊字符、值中存在类似
key=的结构),优先推荐从源头处理数据。
内容的提问来源于stack exchange,提问作者Gorka Lertxundi
相关产品推荐
相关产品推荐

