Spark读取S3中JSON文件到DataFrame后如何还原为原始JSON格式
问题原因
你使用spark.read.json读取JSON文件时,Spark会自动执行JSON解析、Schema推断、数据类型转换操作,还会按照推断的Schema统一字段顺序,原始JSON中的字段排列顺序、可选字段省略规则、数值/字符串的原生表达格式都会在解析阶段被改写。再转pandas的过程又会新增一层类型转换(比如Spark Timestamp类型转pandas datetime类型),自然无法得到和原始完全一致的JSON内容。
解决方案
方案1:读取时保留原始JSON串(可100%还原,推荐)
如果需要同时使用结构化数据做计算、又要保留原始JSON格式,不要直接用spark.read.json解析,先读取为纯文本行存储原始JSON内容,再额外解析结构化字段即可:
from pyspark.sql.functions import from_json, col, struct # 第一步:读取所有JSON文件为纯文本,每行对应一个原始JSON字符串 raw_text_df = spark.read.text("s3://some-bucket/some-key/*.json") # 第二步:按需解析结构化字段,同时保留原始JSON列 # 这里可以直接复用你之前通过spark.read.json得到的df的Schema origin_schema = df.schema processed_df = raw_text_df \ .withColumn("raw_json", col("value")) \ .withColumn("structured_data", from_json(col("value"), origin_schema)) # 需要导出原始JSON时,直接取raw_json列写为文本文件即可 processed_df.select("raw_json").write.mode("overwrite").text("s3://your-target-output-path/")
方案2:仅保留了解析后的结构化DF时的还原方案
如果已经没有原始读取的任务,只能基于现有结构化DF还原,可以使用Spark原生的to_json函数,相比转pandas的还原度更高:
from pyspark.sql.functions import to_json, struct # 可通过options参数调整输出规则,比如是否保留null值字段、时间格式等 json_df = df.select( to_json( struct(*df.columns), options={"ignoreNullFields": "false", "timestampFormat": "yyyy-MM-dd'T'HH:mm:ss.SSSZ"} ).alias("value") ) # 写出为JSON文本文件 json_df.write.mode("overwrite").text("s3://your-target-output-path/")
注意:该方案无法做到100%和原始JSON一致,因为解析阶段丢失的字段顺序、数值原生精度(比如原始JSON中整数被推断为浮点数后导出会带
.0)等信息无法找回。
内容的提问来源于stack exchange,提问作者MykG
相关产品推荐
相关产品推荐

