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

Spark Java实现多列合并为单个JSON结构字段的方法

在Java Spark中合并多列为单个字符串列

方法一:利用struct + to_json生成标准JSON格式

Java Spark中同样支持struct方法,它位于org.apache.spark.sql.functions类下。结合to_json函数可以直接将结构体转换为JSON格式的字符串,这是最简洁的实现方式。

代码示例

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import static org.apache.spark.sql.functions.*;

public class MergeColumns {
    public static void main(String[] args) {
        SparkSession spark = SparkSession.builder()
                .appName("MergeColumnsExample")
                .master("local[*]")
                .getOrCreate();

        // 构造源数据表
        Dataset<Row> sourceDF = spark.createDataFrame(
                spark.sparkContext().parallelize(java.util.Arrays.asList(
                        new Object[]{"Anuj", "Rai", 26}
                )),
                org.apache.spark.sql.types.DataTypes.createStructType(java.util.Arrays.asList(
                        org.apache.spark.sql.types.DataTypes.createStructField("FirstName", org.apache.spark.sql.types.DataTypes.StringType, true),
                        org.apache.spark.sql.types.DataTypes.createStructField("LastName", org.apache.spark.sql.types.DataTypes.StringType, true),
                        org.apache.spark.sql.types.DataTypes.createStructField("Age", org.apache.spark.sql.types.DataTypes.IntegerType, true)
                ))
        );

        // 生成JSON格式的person_details列
        Dataset<Row> targetDF = sourceDF.withColumn(
                "person_details",
                to_json(struct(col("FirstName"), col("LastName"), col("Age")))
        );

        targetDF.show(false);
        spark.stop();
    }
}

输出结果

+---------+--------+---+---------------------------------------+
|FirstName|LastName|Age|person_details                         |
+---------+--------+---+---------------------------------------+
|Anuj     |Rai     |26 |{"FirstName":"Anuj","LastName":"Rai","Age":26}|
+---------+--------+---+---------------------------------------+

方法二:手动拼接成示例指定格式(键不带引号)

如果需要严格匹配你给出的示例格式(键无双引号),可以使用concat结合lit函数手动拼接每个键值对:

代码示例

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import static org.apache.spark.sql.functions.*;

public class MergeColumnsCustom {
    public static void main(String[] args) {
        SparkSession spark = SparkSession.builder()
                .appName("MergeColumnsCustomExample")
                .master("local[*]")
                .getOrCreate();

        Dataset<Row> sourceDF = spark.createDataFrame(
                spark.sparkContext().parallelize(java.util.Arrays.asList(
                        new Object[]{"Anuj", "Rai", 26}
                )),
                org.apache.spark.sql.types.DataTypes.createStructType(java.util.Arrays.asList(
                        org.apache.spark.sql.types.DataTypes.createStructField("FirstName", org.apache.spark.sql.types.DataTypes.StringType, true),
                        org.apache.spark.sql.types.DataTypes.createStructField("LastName", org.apache.spark.sql.types.DataTypes.StringType, true),
                        org.apache.spark.sql.types.DataTypes.createStructField("Age", org.apache.spark.sql.types.DataTypes.IntegerType, true)
                ))
        );

        // 手动拼接生成目标格式的字符串列
        Dataset<Row> targetDF = sourceDF.withColumn(
                "person_details",
                concat(
                        lit("{FirstName:\""), col("FirstName"), lit("\","),
                        lit("LastName:\""), col("LastName"), lit("\","),
                        lit("Age:\""), col("Age").cast("string"), lit("\"}")
                )
        );

        targetDF.show(false);
        spark.stop();
    }
}

输出结果

+---------+--------+---+---------------------------------------+
|FirstName|LastName|Age|person_details                         |
+---------+--------+---+---------------------------------------+
|Anuj     |Rai     |26 |{FirstName:"Anuj",LastName:"Rai",Age:"26"}|
+---------+--------+---+---------------------------------------+

关键说明

  • Java中的struct函数用法和Scala本质一致,通过functions.struct()传入多个列对象即可创建结构体列。
  • 如果需要处理大量列,手动拼接会比较繁琐,可以通过遍历列名的方式动态生成拼接表达式,减少重复代码。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 15:38:15