如何将Java Spark DataFrame转换为指定格式JSON并分区?
Spark DataFrame 转指定格式JSON并控制分组大小实现
实现思路
核心通过分组聚合将每100行数据打包成predictions数组,再包装成目标JSON结构,步骤如下:
- 给每行分配分组ID,确保每组最多包含100条记录;
- 将单行的a、b、c列转换为结构体(对应单个JSON对象);
- 按分组ID聚合,收集组内所有结构体为数组;
- 将数组包装成指定的外层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
相关产品推荐
相关产品推荐

