如何用PySpark窗口函数按ID分组计算首尾行差值并生成新列?
PySpark实现分组首尾col1差值的简洁方案
要实现按ID分区,计算每组最后一行与第一行col1的差值并生成新列col2,最简洁的方式是利用PySpark的窗口函数结合first()和last()聚合函数,具体步骤如下:
1. 导入依赖并初始化环境
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import first, last # 初始化SparkSession spark = SparkSession.builder.appName("GroupFirstLastDiff").getOrCreate()
2. 创建示例数据
# 构造输入数据 data = [ (1, 1), (1, 2), (1, 4), (2, 1), (2, 1), (2, 6), (3, 5), (3, 5), (3, 7) ] df = spark.createDataFrame(data, ["ID", "col1"])
3. 定义窗口规范并计算差值
# 按ID分区的窗口(若需明确行顺序,可添加orderBy,比如orderBy("col1")) window_spec = Window.partitionBy("ID") # 添加col2列:每组最后一行col1减去第一行col1 result_df = df.withColumn( "col2", last("col1", ignorenulls=True).over(window_spec) - first("col1", ignorenulls=True).over(window_spec) )
4. 查看结果
执行result_df.show()后输出如下:
+---+----+----+ | ID|col1|col2| +---+----+----+ | 1| 1| 3| | 1| 2| 3| | 1| 4| 3| | 2| 1| 5| | 2| 1| 5| | 2| 6| 5| | 3| 5| 2| | 3| 5| 2| | 3| 7| 2| +---+----+----+
关键说明
partitionBy("ID")实现按ID分组;last("col1")和first("col1")分别提取每组最后一行和第一行的col1值,ignorenulls=True确保忽略空值(若存在);- 若数据行顺序不确定,建议在窗口规范中添加
orderBy子句(如orderBy("col1"))来明确首尾行的判定逻辑,避免因物理存储顺序导致结果异常。
内容的提问来源于stack exchange,提问作者cnns
相关产品推荐
相关产品推荐

