You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 19:55:24