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

Spark Java中使用Map展开含多指标的行(行转列)

Spark Java 用Map方式实现DataFrame行转列

你需要将包含id、sent、delivered、opened列的DataFrame,转换为每行对应一个指标的长表结构,具体结构如下:

原DataFrame结构:

id | sent  | delivered | opened
--------------------------------
1  | 5     | 3         | 2
2  | 11    | 9         | 4

目标DataFrame结构:

id  | metric_name | metric_value
--------------------------------
1   | sent        |  5
1   | delivered   |  3
1   | opened      |  2
2   | sent        |  11
2   | delivered   |  9
2   | opened      |  4

实现步骤与代码

  1. 定义目标Schema:先构建转换后DataFrame的Schema,包含id(整数类型)、metric_name(字符串类型)、metric_value(整数类型)。
  2. 用flatMap拆分行:通过flatMap将原DataFrame的每一行拆分为多个对应不同指标的行,遍历指标列生成对应Row对象,这就是核心的Map思想落地。
  3. 转换为目标DataFrame:将处理后的RDD转换为符合目标Schema的DataFrame。

具体Java代码如下:

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructField;
import org.apache.spark.sql.types.StructType;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;

public class SparkRowToColumn {
    public static void main(String[] args) {
        // 初始化SparkSession
        SparkSession spark = SparkSession.builder()
                .appName("RowToColumnMap")
                .master("local[*]")
                .getOrCreate();

        // 模拟原DataFrame数据
        List<Row> originalData = Arrays.asList(
                RowFactory.create(1, 5, 3, 2),
                RowFactory.create(2, 11, 9, 4)
        );

        StructType originalSchema = DataTypes.createStructType(new StructField[]{
                DataTypes.createStructField("id", DataTypes.IntegerType, false),
                DataTypes.createStructField("sent", DataTypes.IntegerType, false),
                DataTypes.createStructField("delivered", DataTypes.IntegerType, false),
                DataTypes.createStructField("opened", DataTypes.IntegerType, false)
        });

        Dataset<Row> originalDF = spark.createDataFrame(originalData, originalSchema);
        originalDF.show();

        // 定义需要转换的指标列名
        List<String> metricColumns = Arrays.asList("sent", "delivered", "opened");

        // 构建目标Schema
        StructType targetSchema = DataTypes.createStructType(new StructField[]{
                DataTypes.createStructField("id", DataTypes.IntegerType, false),
                DataTypes.createStructField("metric_name", DataTypes.StringType, false),
                DataTypes.createStructField("metric_value", DataTypes.IntegerType, false)
        });

        // 使用flatMap实现行转列(核心Map逻辑:遍历列生成对应行)
        Dataset<Row> targetDF = spark.createDataFrame(
                originalDF.javaRDD().flatMap(row -> {
                    Integer id = row.getInt(0);
                    List<Row> resultRows = new ArrayList<>();
                    // 遍历每个指标列,映射为新的Row
                    for (String metric : metricColumns) {
                        Integer value = row.getAs(metric);
                        resultRows.add(RowFactory.create(id, metric, value));
                    }
                    return resultRows.iterator();
                }),
                targetSchema
        );

        // 展示结果
        targetDF.show();

        spark.stop();
    }
}

代码说明

  • 选用flatMap而非普通map,是因为一行需要拆分成多行,flatMap可以将单个输入元素映射为多个输出元素。
  • 通过遍历预定义的指标列名列表,从原Row中提取对应值并构建新Row,完全贴合Map方式的转换逻辑。
  • 这种方式灵活性强,只需调整metricColumns列表就能增减需要转换的指标列。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 20:48:45