如何用PySpark/Spark SQL将特定行转为列(替代Join逻辑)
实现需求:将Total行转为列(无需Join逻辑)
可以通过窗口函数实现这个转换,全程不需要使用Join操作,下面提供PySpark API和Spark SQL两种实现方式:
一、PySpark API 实现
1. 创建示例数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("total_row_to_column").getOrCreate() # 构造输入数据 data = [ ("a", "b", "abc", 0), ("a", "b", "total", 9), ("e", "r", "uty", 5), ("e", "r", "total", 44), ("w", "v", "xtr", 3) ] df = spark.createDataFrame(data, ["col1", "col2", "col3", "col4"])
2. 核心转换逻辑
# 定义窗口:按col1、col2分组,确保同组数据能共享total值 window_spec = Window.partitionBy("col1", "col2") # 提取同组内total行的col4值作为新列,过滤掉原total行后填充空值为0 df_result = df.withColumn( "total", F.max(F.when(F.col("col3") == "total", F.col("col4"))).over(window_spec) ).filter(F.col("col3") != "total") \ .fillna(0, subset=["total"]) # 查看结果 df_result.show()
二、Spark SQL 实现
1. 创建临时视图并执行SQL
# 先将DataFrame注册为临时视图 df.createOrReplaceTempView("source_table") # 执行SQL查询 spark.sql(""" SELECT col1, col2, col3, col4, COALESCE(MAX(CASE WHEN col3 = 'total' THEN col4 END) OVER(PARTITION BY col1, col2), 0) AS total FROM source_table WHERE col3 != 'total' """).show()
逻辑说明
核心思路是利用窗口函数将同一col1、col2组合的行划分为一个窗口,在窗口内提取col3='total'对应的col4值;因为每个分组最多只有一条total行,用MAX或FIRST函数都能准确提取目标值;最后过滤掉原始的total行,并用COALESCE或fillna将无total值的分组填充为0,完全规避了Join操作的开销。
内容的提问来源于stack exchange,提问作者rangerover
相关产品推荐
相关产品推荐

