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

如何在Apache Spark中为Dataset/DataFrame添加JSON对象并生成自定义JSON

使用Apache Spark Dataset API构建自定义嵌套JSON结构

刚好我之前处理过类似的需求,完全可以用Spark的内置函数来实现你要的效果——把alerts Dataset作为inventory里名为ALERT的JSON对象(或数组)嵌套进去。下面分两种常见场景给你具体的代码示例:

场景1:一对一关联(一个inventory对应一个alert)

假设两个Dataset有共同的关联键(比如item_id),我们先关联数据,再把alerts的字段打包成ALERT嵌套对象:

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import static org.apache.spark.sql.functions.struct;

// 加载你的inventory Dataset(沿用你提供的加载逻辑)
Dataset<Row> inventory = spark.read()
    .option("multiLine", true)
    .option("mode", "PERMISSIVE")
    .json("C:\\Users\\phyadavi\\LearningAn...");

// 加载alerts Dataset(替换为你的实际路径)
Dataset<Row> alerts = spark.read().json("path/to/your/alerts.json");

// 用关联键(比如item_id)关联两个Dataset,left_outer确保inventory的所有数据都保留
Dataset<Row> joinedData = inventory.join(alerts, inventory.col("item_id").equalTo(alerts.col("item_id")), "left_outer");

// 将alerts的所有字段打包成名为ALERT的结构体,输出JSON时会自动转为嵌套对象
Dataset<Row> finalDataset = joinedData
    .select(
        inventory.col("*"), // 保留inventory原有的所有字段
        struct(alerts.col("*")).alias("ALERT") // 把alerts的字段组合成ALERT嵌套对象
    );

// 输出结果,multiLine=true会让每条数据格式化输出
finalDataset.write()
    .option("multiLine", true)
    .mode(SaveMode.Overwrite)
    .json("path/to/your/output.json");

如果需要把ALERT字段输出为JSON字符串(而不是嵌套的JSON对象),可以用to_json函数转换:

import static org.apache.spark.sql.functions.to_json;

Dataset<Row> finalDatasetWithJsonString = joinedData
    .select(
        inventory.col("*"),
        to_json(struct(alerts.col("*"))).alias("ALERT")
    );

场景2:一对多关联(一个inventory对应多个alerts)

如果一个inventory条目对应多条alert记录,我们需要先把alerts按关联键分组,收集成数组后再关联:

import static org.apache.spark.sql.functions.collect_list;

// 先把alerts按item_id分组,收集所有alert为数组
Dataset<Row> groupedAlerts = alerts
    .groupBy("item_id")
    .agg(collect_list(struct(alerts.col("*"))).alias("ALERT"));

// 关联到inventory,此时ALERT字段是一个包含多个alert对象的数组
Dataset<Row> finalDataset = inventory.join(groupedAlerts, "item_id", "left_outer");

关键函数说明

  • struct():把多个字段组合成一个Spark结构体类型,输出JSON时会自动转为嵌套的JSON对象
  • to_json():将结构体或复杂类型转换为JSON格式的字符串
  • collect_list():分组时将多行数据收集为数组,适合处理一对多的关联场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:15:06