Spark SQL及PySpark环境下如何实现列转行(宽表转长表)操作
宽表转KV结构长表(逆透视)实现方案
你需要的是逆透视(Unpivot)操作,和pivot(行转列)是反向逻辑,以下分别给出PySQL(Spark SQL)和Spark Python DataFrame两种场景的高效实现:
1. Spark SQL 实现(基于stack函数)
无需手动拼接union语句,直接使用内置stack函数即可完成转换,示例代码如下:
SELECT ID, ValueVv, ValueDesc FROM 你的原表名 LATERAL VIEW stack(3, 'Value1', Value1, 'Value2', Value2, 'Value40', Value40) tmp AS ValueDesc, ValueVv
参数说明:
- stack第一个参数
3代表要把多少列转成行,此处对应Value1/Value2/Value40共3个列 - 后续每两个参数为一组:第一个是转成ValueDesc的列名字面量,第二个是对应列的取值
- 最后别名
ValueDesc, ValueVv对应输出的两个字段名
如果需要转换的列数量多,可动态拼接SQL字段部分,无需手动逐个书写。
2. Spark Python DataFrame 实现
提供三种适配不同版本的实现方式,性能均优于手动union:
通用兼容版(所有Spark版本可用)
基于selectExpr拼接stack表达式实现:
# 1. 定义需要转换的列,可自动匹配前缀批量生成 value_cols = [c for c in df.columns if c.startswith("Value")] # 2. 拼接stack函数参数 stack_params = [f"'{col}', {col}" for col in value_cols] stack_expr = f"stack({len(value_cols)}, {','.join(stack_params)}) as (ValueDesc, ValueVv)" # 3. 执行转换 result_df = df.selectExpr("ID", stack_expr)
Spark 3.4+ 原生unpivot实现
3.4及以上版本Spark提供了原生unpivot方法,语法更直观:
value_cols = [c for c in df.columns if c.startswith("Value")] result_df = df.unpivot( ids=["ID"], values=value_cols, variableColumnName="ValueDesc", valueColumnName="ValueVv" )
Pandas风格melt实现(基于Pandas API on Spark)
如果你熟悉pandas的melt用法,可直接用Pandas API on Spark实现相同逻辑:
import pyspark.pandas as ps # Spark DataFrame转Pandas on Spark DataFrame ps_df = df.to_pandas_on_spark() # 执行melt转换 result_ps_df = ps_df.melt( id_vars=["ID"], value_vars=value_cols, var_name="ValueDesc", value_name="ValueVv" ) # 转回标准Spark DataFrame result_df = result_ps_df.to_spark()
如果需要过滤转换后的空值,在转换完成后加where ValueVv is not null条件即可。
内容的提问来源于stack exchange,提问作者Adam
相关产品推荐
相关产品推荐

