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
相关产品推荐
相关产品推荐

