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

PySpark如何基于其他列对DataFrame指定列执行移位操作

PySpark 实现指定规则的Col4左移计算

核心规则对齐

你需要的移位逻辑可以直接通过窗口函数实现,计算边界如下:

  • 计算隔离:所有移位计算按Col1分区,不同Col1值的数据互不影响
  • 排序依据:每个分区内先按Col3升序(区分Col2的不同连续序列批次),同一Col3批次内按Col2升序,保证行顺序匹配连续值变化的逻辑
  • 取值规则:取排序后上一行的Col4作为当前行的shift_col4,每个分区的第一行因为没有前置行,取值为null

实现代码

先导入依赖:

from pyspark.sql import Window
import pyspark.sql.functions as F

定义窗口规则并计算新列:

# 定义窗口分区与排序规则
calc_window = Window.partitionBy("Col1").orderBy("Col3", "Col2")

# 生成移位结果列
result_df = source_df.withColumn("shift_col4", F.lag("Col4", 1).over(calc_window))

效果验证

代码运行后输出结果和你给出的示例完全匹配:

ID  Col1 Col2 Col3 Col4 shift_col4
1    1   10   1    4     null
2    1   11   1    8     4
3    1   12   1   12     8
4    1   1    2   16     12
5    1   3    2   20     16
4    2   1    1   16     null
5    2   4    1   20     16

如果你的源数据ID字段顺序和Col3、Col2排序后的顺序完全一致,排序字段也可以换成ID,计算性能一致;优先用Col3、Col2排序的写法兼容性更好,不会因为ID乱序导致连续序列判断错误。

内容的提问来源于stack exchange,提问作者Mithun

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 17:09:20