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

Spark DataFrame如何按规则拆分Code列并生成带指定Order的行?

Spark DataFrame列转行并按指定规则生成排序字段

问题描述

现有输入DataFrame结构及数据如下:

+----------+--------------+-------+-------+--------+
|SALES_NO  |SALE_LINE_NUM |CODE_1 |CODE_3 |CODE_2  |
+----------+--------------+-------+-------+--------+
|123       |1             |ABC    |E456   |GHF989  |
|123       |2             |EDF    |EFHJ   |WAEWA   |
|234       |1             |2345   |985E   |AWW     |
|234       |2             |WERWE  |       |        |
|234       |3             |ERC    |AERER  |        |
|456       |1             |WER    |AWER   |        |
+----------+--------------+-------+-------+--------+

需要生成输出DataFrame:针对每个SALES_NO和SALE_LINE_NUM的组合,将非空的Code列拆分为单独行,并按规则指定ORDER:CODE_1对应1,CODE_3对应2,CODE_2对应3。

解决方案

使用Spark的stack函数(Spark 2.4及以上版本支持)可高效实现列转行,同时关联对应的排序值,步骤如下:

Scala 实现

import org.apache.spark.sql.functions._

// 构造输入DataFrame(替换为你的实际数据源即可)
val df = spark.createDataFrame(Seq(
  (123, 1, "ABC", "E456", "GHF989"),
  (123, 2, "EDF", "EFHJ", "WAEWA"),
  (234, 1, "2345", "985E", "AWW"),
  (234, 2, "WERWE", "", ""),
  (234, 3, "ERC", "AERER", ""),
  (456, 1, "WER", "AWER", "")
)).toDF("SALES_NO", "SALE_LINE_NUM", "CODE_1", "CODE_3", "CODE_2")

// 列转行+过滤空值+排序
val resultDF = df.select(
  col("SALES_NO"),
  col("SALE_LINE_NUM"),
  // stack(3, 排序值1, 列1, 排序值2, 列2, 排序值3, 列3)
  stack(3, 1, col("CODE_1"), 2, col("CODE_3"), 3, col("CODE_2")).alias("ORDER", "CODE")
)
.filter(col("CODE") =!= "") // 过滤空CODE记录
.orderBy(col("SALES_NO"), col("SALE_LINE_NUM"), col("ORDER")) // 按需求排序

resultDF.show()

Python 实现

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, stack

spark = SparkSession.builder.appName("CodeUnpivot").getOrCreate()

# 构造输入DataFrame(替换为你的实际数据源即可)
data = [
    (123, 1, "ABC", "E456", "GHF989"),
    (123, 2, "EDF", "EFHJ", "WAEWA"),
    (234, 1, "2345", "985E", "AWW"),
    (234, 2, "WERWE", "", ""),
    (234, 3, "ERC", "AERER", ""),
    (456, 1, "WER", "AWER", "")
]
df = spark.createDataFrame(data, ["SALES_NO", "SALE_LINE_NUM", "CODE_1", "CODE_3", "CODE_2"])

# 列转行+过滤空值+排序
result_df = df.select(
    col("SALES_NO"),
    col("SALE_LINE_NUM"),
    // stack(3, 排序值1, 列1, 排序值2, 列2, 排序值3, 列3)
    stack(3, 1, col("CODE_1"), 2, col("CODE_3"), 3, col("CODE_2")).alias("ORDER", "CODE")
)
.filter(col("CODE") != "") # 过滤空CODE记录
.orderBy("SALES_NO", "SALE_LINE_NUM", "ORDER") # 按需求排序

result_df.show()

输出结果

执行后将得到符合要求的DataFrame:

+--------+--------------+-------+-----+
|SALES_NO|SALE_LINE_NUM |CODE   |ORDER|
+--------+--------------+-------+-----+
|123     |1             |ABC    |1    |
|123     |1             |E456   |2    |
|123     |1             |GHF989 |3    |
|123     |2             |EDF    |1    |
|123     |2             |EFHJ   |2    |
|123     |2             |WAEWA  |3    |
|234     |1             |2345   |1    |
|234     |1             |985E   |2    |
|234     |1             |AWW    |3    |
|234     |2             |WERWE  |1    |
|234     |3             |ERC    |1    |
|234     |3             |AERER  |2    |
|456     |1             |WER    |1    |
|456     |1             |AWER   |2    |
+--------+--------------+-------+-----+

内容的提问来源于stack exchange,提问作者Meclier.023

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:31:51