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
实现步骤与代码
- 定义目标Schema:先构建转换后DataFrame的Schema,包含
id(整数类型)、metric_name(字符串类型)、metric_value(整数类型)。 - 用flatMap拆分行:通过
flatMap将原DataFrame的每一行拆分为多个对应不同指标的行,遍历指标列生成对应Row对象,这就是核心的Map思想落地。 - 转换为目标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
相关产品推荐
相关产品推荐

