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

如何将Java Spark DataFrame转换为指定格式JSON并分区?

Spark DataFrame 转指定格式JSON并控制分组大小实现

实现思路

核心通过分组聚合将每100行数据打包成predictions数组,再包装成目标JSON结构,步骤如下:

  1. 给每行分配分组ID,确保每组最多包含100条记录;
  2. 将单行的a、b、c列转换为结构体(对应单个JSON对象);
  3. 按分组ID聚合,收集组内所有结构体为数组;
  4. 将数组包装成指定的外层JSON格式。

Python 代码实现

假设原始DataFrame名为df:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 1. 生成分组ID:每100行一组
window_spec = Window.orderBy(F.monotonically_increasing_id())
df_with_group = df.withColumn(
    "group_id",
    F.floor(F.row_number().over(window_spec) / 100)
)

# 2. 将单行数据转为结构体(对应单个JSON对象)
df_struct = df_with_group.withColumn(
    "row_json",
    F.struct(F.col("a"), F.col("b"), F.col("c"))
)

# 3. 按分组聚合,收集结构体为predictions数组
df_agg = df_struct.groupBy("group_id").agg(
    F.collect_list("row_json").alias("predictions")
)

# 4. 生成目标格式的JSON字符串
df_result = df_agg.withColumn(
    "final_json",
    F.to_json(F.struct(F.col("predictions")))
)

# 可选:输出结果(每个文件对应一个分组的JSON)
df_result.drop("group_id", "predictions").write.mode("overwrite").text("your_output_path")

Scala 代码实现

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

val df: DataFrame = // 你的原始DataFrame

// 1. 生成分组ID
val windowSpec = Window.orderBy(monotonically_increasing_id())
val dfWithGroup = df.withColumn(
    "group_id",
    floor(row_number().over(windowSpec) / 100)
)

// 2. 单行转结构体
val dfStruct = dfWithGroup.withColumn(
    "row_json",
    struct(col("a"), col("b"), col("c"))
)

// 3. 聚合生成predictions数组
val dfAgg = dfStruct.groupBy("group_id").agg(
    collect_list("row_json").alias("predictions")
)

// 4. 生成目标JSON
val dfResult = dfAgg.withColumn(
    "final_json",
    to_json(struct(col("predictions")))
)

// 可选:保存结果
dfResult.drop("group_id", "predictions").write.mode("overwrite").text("your_output_path")

关键说明

  • 分组逻辑:用row_number()生成行号后整除100,确保每个分组严格不超过100条数据;若无需严格顺序,也可使用hash函数结合分区数实现,但行号方式可控性更强。
  • JSON格式:Spark的to_json会生成标准带双引号的JSON(符合JSON规范),若需非标准格式,可手动拼接字符串,但不推荐。
  • 输出:用text格式保存时,每个分组的JSON会单独成为一行,方便后续读取处理。

内容的提问来源于stack exchange,提问作者Joe C

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 07:30:54