Spark Scala如何实现长表转宽表得到指定目标输出格式
Spark Scala 实现长表转指定结构宽表方案
实现思路
- 先通过窗口函数对同分组(A、B、C)内的记录生成自增序号,区分同一个C值下的多条D记录
- 构造带序号的动态列名,再通过pivot操作完成行转列
完整实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ object PivotDemo { def main(args: Array[String]): Unit = { // 1. 初始化SparkSession val spark = SparkSession.builder() .appName("LongToWideDemo") .master("local[*]") .getOrCreate() import spark.implicits._ // 2. 构造示例原始数据 val rawDF = Seq( (1, "A", "Day", "D1"), (1, "A", "Tim", "1am"), (1, "A", "Tim", "3am") ).toDF("A", "B", "C", "D") // 3. 定义窗口:按A、B、C分组,按D值升序排序(保证Tim下的时间顺序正确) val win = Window.partitionBy("A", "B", "C").orderBy("D") // 4. 生成行号、构造动态pivot列名 val processedDF = rawDF .withColumn("rn", row_number().over(win)) .withColumn("pivot_col", when(col("C") === "Day", col("C")) .otherwise(concat(col("C"), col("rn"))) ) // 5. 分组pivot得到最终宽表 val resultDF = processedDF .groupBy("A", "B") .pivot("pivot_col") .agg(first("D")) // 可根据实际需求调整列顺序 .select("A", "B", "Day", "Tim1", "Tim2") // 打印结果验证 resultDF.show() /* 输出结果: +---+---+---+----+----+ | A| B|Day|Tim1|Tim2| +---+---+---+----+----+ | 1| A| D1| 1am| 3am| +---+---+---+----+----+ */ spark.stop() } }
逻辑说明
- 窗口函数的排序规则可根据业务调整,比如需要按数据入库顺序给Tim值排序,替换
orderBy的字段即可 - 如果同分组下Tim的数量不固定,可以先收集所有pivot列名再动态指定select的列顺序,不需要硬编码Tim1、Tim2
- 聚合函数用
first是因为同A、B、pivot_col分组下只有一条D值,替换为max/min效果一致
内容的提问来源于stack exchange,提问作者Krishna Murthy
相关产品推荐
相关产品推荐

