PySpark使用collect_list聚合时如何去除列值中的转义字符
Spark CSV转分组JSON避免转义字符的解决方案
问题根源
你提前将结构化JSON对象序列化为字符串后再存入DataFrame,Spark写入JSON文件时会自动对字符串内部的双引号做转义处理,所以才会出现大量反斜杠。之前去掉json.dumps后出现等号格式,是因为myObj.as_json()返回的是对象的默认字符串表示而非字典结构,Spark会直接把该字符串作为resource字段的值写入。
修复步骤
1. 调整RDD转换逻辑,保留结构化对象而非提前转字符串
确保mapToJson返回的是Python字典结构,不要提前做JSON序列化:
import json from pyspark.sql import Row def mapToJson(instance): myObj = MyCustomObject() # 确保resource字段是字典结构:如果as_json返回的是JSON字符串,就用json.loads转成字典 obj_json = myObj.as_json() _json = json.loads(obj_json) if isinstance(obj_json, str) else obj_json return { "fullUrl": "urn:uuid", "resource": _json } df1 = spark.read.format("csv").option("delimiter","|").option("header","true").load(filepathsrc) # 直接返回字典,Spark会自动识别为结构化的StructType类型 rddJson = df1.rdd.map(lambda instance: mapToJson(instance)) df = spark.createDataFrame(rddJson)
2. 调整分组聚合逻辑
直接对结构化字段做collect_list,不需要额外做类型转换:
from pyspark.sql import Window from pyspark.sql.functions import row_number, lit, col, collect_list, struct w = Window().orderBy(lit('A')) df = df.withColumn("ROW_ID", row_number().over(w)) \ .withColumn("group_num",(row_number().over(Window.orderBy("ROW_ID"))-1) % 2 ) \ .groupby(col("group_num")) \ # 直接收集结构化对象到entry数组 .agg(collect_list(struct("fullUrl", "resource")).alias("entry")) \ .drop("group_num") \ .withColumn("resourceType", lit("Bundle"))
3. 写入JSON文件
去掉无效的quote配置,直接写入即可:
df.write.format("json")\ .mode('overwrite')\ .option("maxRecordsPerFile",1)\ .save(filepathtgt)
最终效果
写入的JSON文件会完全符合你的预期,entry字段是JSON对象数组,没有多余的转义字符。
内容的提问来源于stack exchange,提问作者aiman
相关产品推荐
相关产品推荐

