PySpark:将每个嵌套JSON存入DataFrame单个单元格时多文件重复问题
解决多文件处理时DataFrame中data列重复同一JSON的问题
问题根源
你之前用df.withColumn("data", lit(df.toJSON().first()))的方式,在合并多文件的DataFrame后调用first(),只会取合并后DataFrame的第一条记录对应的JSON,导致所有行的data列都复用这个值,完全不符合每个文件对应自身完整JSON的需求。
正确处理思路
核心是单个文件单独处理,每个文件对应DataFrame中的一行,分别生成DataType和该文件的完整JSON内容,再合并所有文件的结果后分区写入。
方案1:读取原始JSON文件内容(推荐)
如果需要保留文件原始的JSON格式(避免Spark解析后序列化带来的格式/顺序变化),直接读取文件的完整文本内容:
from pyspark.sql import SparkSession from pyspark.sql.functions import lit import os spark = SparkSession.builder.appName("JSONToParquet").getOrCreate() # 替换为你的JSON文件存储路径 json_dir = "s3://your-bucket/json-files/" json_files = [f"{json_dir}{f}" for f in os.listdir(json_dir) if f.endswith(".json")] total_df = None for file_path in json_files: # 读取整个文件的原始文本内容 raw_json_df = spark.read.text(file_path, wholetext=True) # 从文件名提取DataType(根据你的文件名规则调整逻辑) file_name = os.path.basename(file_path) data_type = file_name.split("_")[0] # 示例:假设文件名是user_data_001.json,取user_data # 添加DataType列并重命名原始文本列为data processed_df = raw_json_df.withColumn("DataType", lit(data_type))\ .withColumnRenamed("value", "data") # 合并到总DataFrame if total_df is None: total_df = processed_df else: total_df = total_df.union(processed_df) # 按DataType分区写入Parquet(替换为你的输出路径) total_df.write.partitionBy("DataType").parquet("s3://your-bucket/output-parquet/")
方案2:基于解析后的DataFrame生成JSON
如果需要使用Spark解析后的JSON结构(比如统一格式),可以将单个文件的DataFrame序列化为JSON字符串:
from pyspark.sql import SparkSession from pyspark.sql.functions import lit import os spark = SparkSession.builder.appName("JSONToParquet").getOrCreate() json_dir = "s3://your-bucket/json-files/" json_files = [f"{json_dir}{f}" for f in os.listdir(json_dir) if f.endswith(".json")] total_df = None for file_path in json_files: # 读取并解析JSON文件 parsed_df = spark.read.json(file_path) # 提取DataType file_name = os.path.basename(file_path) data_type = file_name.split("_")[0] # 将当前文件的所有记录序列化为JSON字符串(数组格式) json_records = parsed_df.toJSON().collect() full_json_str = "[" + ",".join(json_records) + "]" if len(json_records) > 1 else json_records[0] # 创建单条记录的DataFrame single_row_df = spark.createDataFrame([(data_type, full_json_str)], schema=["DataType", "data"]) # 合并到总DataFrame if total_df is None: total_df = single_row_df else: total_df = total_df.union(single_row_df) # 分区写入Parquet total_df.write.partitionBy("DataType").parquet("s3://your-bucket/output-parquet/")
关键注意事项
- 绝对不要在合并后的大DataFrame中使用
first()获取JSON,这只会取第一条记录的内容 - 单个文件处理能保证每个
data列的值严格对应当前文件的完整JSON - 若需保留原始文件格式,优先用方案1的
wholetext=True读取方式
内容的提问来源于stack exchange,提问作者ujp
相关产品推荐
相关产品推荐

