PySpark如何将DataFrame每N条记录分组输出为指定结构的JSON文件
PySpark实现CSV分组转带外层结构JSON方案
核心思路
先给每500条记录分配同一个分组ID,按分组ID聚合将同组记录收集为列表,再追加时间字段后输出,每个分组对应一条符合要求的JSON结构。
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import row_number, collect_list, struct, lit, col from pyspark.sql.window import Window from datetime import datetime, timezone # 初始化SparkSession(你已有可跳过) spark = SparkSession.builder.appName("csv2json").getOrCreate() # 读取CSV(你已有可跳过) df = spark.read.csv("你的CSV文件路径", header=True, inferSchema=True) # 1. 分配分组ID,每500条为一组,按row_id排序 window_spec = Window.orderBy("row_id") df_with_group = df.withColumn("group_id", (row_number().over(window_spec) - 1) // 500) # 2. 按分组聚合,构造外层结构 # 将每行记录转为结构体,方便收集为列表 df_record = df_with_group.select("group_id", struct(*df.columns).alias("record")) # 聚合+追加更新时间 df_output = df_record.groupBy("group_id") \ .agg(collect_list("record").alias("entry")) \ .withColumn("last_updated", lit(datetime.now(timezone.utc).isoformat().replace("+00:00", "Z"))) \ .drop("group_id") # 3. 输出JSON,每个文件1条记录(即每个文件对应1个你要求的完整JSON结构) df_output.write \ .mode("overwrite") \ .option("maxRecordsPerFile", 1) \ .json("你的输出目录路径")
可选优化(针对超大数据量)
如果不需要严格按row_id排序分组,可以用monotonically_increasing_id()生成临时ID,避免全局排序的性能损耗:
# 替换上述分配分组ID的步骤即可 from pyspark.sql.functions import monotonically_increasing_id df_with_group = df.withColumn("temp_id", monotonically_increasing_id()) \ .withColumn("group_id", (col("temp_id") - 1) // 500) \ .drop("temp_id")
效果说明
- 输出的每个JSON文件内就是你要求的完整结构,
entry字段包含500条原CSV记录,last_updated为当前UTC时间,格式和你示例完全一致 - 500万条记录总共会生成10000个JSON文件,和需求匹配
内容的提问来源于stack exchange,提问作者aiman
相关产品推荐
相关产品推荐

