PySpark使用窗口函数判断DataFrame列值是否逐行递增(无需lag)
解决方案
核心思路:无需使用lag函数,通过带滑动行范围的窗口+collect_list实现批量多列的递增校验,自动适配任意数量的待校验数值列,无需提前获知列名。
前置说明
- 窗口定义:按
id分组,按datum升序排序(和样例数据的行顺序保持一致),窗口行范围限定为「上1行 ~ 当前行」 - 自动识别待校验列:排除
id/datum/lfd三个固定维度列,剩余所有数值列均自动纳入校验范围,生成对应的FLAG_INCREASE_<列名>标记列
完整可运行代码
import pyspark.sql.functions as F from pyspark.sql.window import Window from pyspark.sql import Row # 样例数据初始化 row = Row("id", "datum", "lfd", "some_value", "some_value2") df = spark.sparkContext.parallelize([ row(1, "2015-01-01", 4, 20.0, 20.0), row(1, "2015-01-06", 3, 10.0, 20.0), row(1, "2015-01-07", 2, 25.0, 20.0), row(1, "2015-01-12", 1, 30.0, 20.0), row(2, "2015-01-01", 4, 5.0, 20.0), row(2, "2015-01-06", 3, 30.0, 20.0), row(2, "2015-01-12", 1, 20.0, 20.0) ]).toDF().withColumn("datum", F.col("datum").cast("date")) # 1. 定义窗口 w = Window.partitionBy("id").orderBy("datum").rowsBetween(-1, 0) # 2. 自动识别待校验列(排除固定维度列) fixed_cols = {"id", "datum", "lfd"} check_cols = [c for c in df.columns if c not in fixed_cols] # 3. 批量生成所有校验列的逻辑 flag_cols = [ F.when(F.size(F.collect_list(F.col(c)).over(w)) == 1, 0) .when(F.collect_list(F.col(c)).over(w)[1] > F.collect_list(F.col(c)).over(w)[0], 1) .otherwise(0) .alias(f"FLAG_INCREASE_{c}") for c in check_cols ] # 4. 生成最终结果 result_df = df.select("*", *flag_cols) result_df.show()
运行输出结果
+---+----------+---+----------+-----------+------------------------+-------------------------+ | id| datum|lfd|some_value|some_value2|FLAG_INCREASE_some_value|FLAG_INCREASE_some_value2| +---+----------+---+----------+-----------+------------------------+-------------------------+ | 1|2015-01-01| 4| 20.0| 20.0| 0| 0| | 1|2015-01-06| 3| 10.0| 20.0| 0| 0| | 1|2015-01-07| 2| 25.0| 20.0| 1| 0| | 1|2015-01-12| 1| 30.0| 20.0| 1| 0| | 2|2015-01-01| 4| 5.0| 20.0| 0| 0| | 2|2015-01-06| 3| 30.0| 20.0| 1| 0| | 2|2015-01-12| 1| 20.0| 20.0| 0| 0| +---+----------+---+----------+-----------+------------------------+-------------------------+
方案说明
- 完全未使用lag函数,符合需求
- 自动适配任意数量的待校验列,新增/删除校验列无需修改代码逻辑
- 标记规则完全符合要求:分组内首行固定为0,当前行值大于上一行时为1,其余情况为0
内容的提问来源于stack exchange,提问作者rkraft
相关产品推荐
相关产品推荐

