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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 00:15:05