如何使用PySpark将多列转换为「产品」与「状态」两列?
PySpark实现宽表转长表(列转行)方法
原始数据格式
Date, product_a, product_b, product_c 2021-01-01, True , True, False 2021-01-02, False, True, False
目标数据格式
Date, product, status 2021-01-01, a, True 2021-01-01, b, True 2021-01-01, c, False 2021-01-02, a, False 2021-01-02, b, True 2021-01-02, c, False
实现方案
方法一:手动指定列转换(适合固定数量的产品列)
使用PySpark的stack函数直接实现列转行,代码如下:
from pyspark.sql import SparkSession # 初始化Spark会话 spark = SparkSession.builder.appName("wide_to_long_conversion").getOrCreate() # 构造原始DataFrame raw_data = [ ("2021-01-01", True, True, False), ("2021-01-02", False, True, False) ] raw_columns = ["Date", "product_a", "product_b", "product_c"] df = spark.createDataFrame(raw_data, raw_columns) # 执行转换 transformed_df = df.selectExpr( "Date", "stack(3, 'a', product_a, 'b', product_b, 'c', product_c) as (product, status)" ) # 查看结果 transformed_df.show()
stack(n, ...)参数说明:n是需要转换的列总数;后续成对传入产品别名(如'a'对应product_a)和列名,最后指定输出的列名(product, status)。
方法二:动态生成转换逻辑(适合数量不固定的产品列)
如果产品列数量不确定,可以通过遍历列名动态生成stack表达式,避免手动编写:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("dynamic_wide_to_long").getOrCreate() # 构造原始DataFrame(同方法一) raw_data = [ ("2021-01-01", True, True, False), ("2021-01-02", False, True, False) ] raw_columns = ["Date", "product_a", "product_b", "product_c"] df = spark.createDataFrame(raw_data, raw_columns) # 筛选出所有产品列(排除Date列) product_columns = [col for col in df.columns if col.startswith("product_")] # 动态生成stack函数的参数 stack_args = [] for col in product_columns: # 提取产品别名(去掉product_前缀) product_alias = col.replace("product_", "") stack_args.extend([f"'{product_alias}'", col]) # 构建完整的stack表达式 stack_expression = f"stack({len(product_columns)}, {', '.join(stack_args)}) as (product, status)" # 执行转换 transformed_df = df.selectExpr("Date", stack_expression) # 查看结果 transformed_df.show()
内容的提问来源于stack exchange,提问作者Jerry Yang
相关产品推荐
相关产品推荐

