如何在Apache Spark(Java)中合并两个Dataset生成嵌套JSON对象列表
利用Spark原生特性实现嵌套JSON生成(Java版)
问题背景
你是Apache Spark(Java)新手,需要将两个Dataset合并生成包含多个嵌套JSON对象的文本文件,现有字符串拼接的方案,希望用Spark原生API优化。两个数据集结构如下:
数据集1:Dataset<Row> firstToSecondGeneration
name|ch1|ch2 |ch3 |ch99 Bob|Joe|James| | Sue|Joe|James| | John| | | |Johnny
数据集2:Dataset<Row> secondToThirdGeneration
chName| gChname Joe| Joe Jr. Joe| Josephine James| James Jr. James| Jamie Johnny| Johnny Jr.
期望生成每行一个嵌套JSON对象的文本文件,结构如你提供的嵌套格式。
解决方案步骤
以下是完全基于Spark原生函数的实现,避免字符串拼接,充分利用Spark分布式处理能力:
1. 处理第一代到第二代的数据集,提取非空子节点
首先将所有ch1到ch99列转为数组,过滤掉空值后,展开为每个name对应单个childName的行:
import org.apache.spark.sql.*; import org.apache.spark.sql.functions.*; import java.util.ArrayList; import java.util.List; // 构造所有ch列的名称列表 List<String> chColumns = new ArrayList<>(); for (int i = 1; i <= 99; i++) { chColumns.add("ch" + i); } // 转换并过滤空child Dataset<Row> childDs = firstToSecondGeneration // 将所有ch列合并为数组 .withColumn("childNames", array(chColumns.stream().map(col).toArray(Column[]::new))) // 过滤数组中的空字符串 .withColumn("childNames", filter(col("childNames"), c -> not(c.equalTo(lit(""))))) // 展开数组为每行一个childName .select(col("name"), explode(col("childNames")).alias("childName"));
2. 关联第二代到第三代的数据集,获取孙节点信息
通过childName与chName关联,匹配每个子节点对应的孙节点:
Dataset<Row> childGrandChildDs = childDs .join(secondToThirdGeneration, childDs.col("childName").equalTo(secondToThirdGeneration.col("chName")), "inner") .select(col("name"), col("childName"), col("gChname").alias("grandChildName"));
3. 聚合孙节点,生成子节点的嵌套结构
按name和childName分组,将孙节点聚合成数组并包装为grandChildren结构:
Dataset<Row> childWithGrandChildrenDs = childGrandChildDs .groupBy("name", "childName") // 将每个grandChildName包装为对象,再聚合成数组 .agg(collect_list(struct(col("grandChildName"))).alias("grandChildren"));
4. 聚合子节点,生成最终的嵌套结构
按name分组,将子节点及其孙节点结构聚合成children数组,得到目标数据集:
Dataset<Row> finalDs = childWithGrandChildrenDs .groupBy("name") // 将每个childName和grandChildren包装为对象,再聚合成数组 .agg(collect_list(struct(col("childName"), col("grandChildren"))).alias("children"));
5. 写入JSON文件
Spark的json写入默认采用JSON Lines格式(每行一个JSON对象),正好符合需求:
finalDs.write() .mode(SaveMode.Overwrite) // 根据需求选择写入模式 .json("/your/output/path");
关键特性说明
- 使用
array+filter批量处理多列的空值,避免逐个列判断的冗余代码 - 用
struct和collect_list实现嵌套JSON结构的生成,完全依赖Spark原生API,性能更优且避免字符串拼接的出错风险 - 全程基于分布式Dataset操作,适合大规模数据处理
内容的提问来源于stack exchange,提问作者saugust
相关产品推荐
相关产品推荐

