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
相关产品推荐
相关产品推荐

