You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 15:18:00