PySpark技术实现:如何过滤掉特定条件之前的数据行
在PySpark中移除DataFrame特定条件之前的所有数据行
需求说明
需要移除所有Time早于LatestRecord = 1对应Time的数据行,保留该时间点及之后的所有记录。
解决方案
核心思路是先定位每个ID分组中符合条件的目标时间,再基于该时间过滤数据:
- 导入依赖模块
from pyspark.sql import functions as F from pyspark.sql.window import Window
- 添加目标时间列
通过窗口函数,将每个ID分组内LatestRecord = 1对应的Time广播到该分组的所有行:
window_spec = Window.partitionBy("ID") df_with_target = df.withColumn( "target_time", # 取分组内符合条件的Time,max确保即使有重复值也能正确获取目标时间 F.max(F.when(F.col("LatestRecord") == 1, F.col("Time"))).over(window_spec) )
- 过滤并清理数据
保留Time大于等于目标时间的行,最后移除临时生成的target_time列:
result_df = df_with_target.filter(F.col("Time") >= F.col("target_time")).drop("target_time")
完整可运行代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化Spark环境 spark = SparkSession.builder.appName("FilterBeforeCondition").getOrCreate() sc = spark.sparkContext # 创建示例DataFrame df = spark.createDataFrame(sc.parallelize([ ['A', '2019-01-01', None, None, None], ['A', '2019-01-02', None, None, None], ['A', '2019-01-03', 'O', 'O', 1], ['A', '2019-01-04', 'O', 'P', 2], ['A', '2019-01-05', 'O', 'P', 3], ['A', '2019-01-06', None, None, None], ['A', '2019-01-07', None, None, 4], ]), ['ID', 'Time', 'State', 'State2', 'LatestRecord']) # 执行过滤逻辑 window_spec = Window.partitionBy("ID") df_with_target = df.withColumn( "target_time", F.max(F.when(F.col("LatestRecord") == 1, F.col("Time"))).over(window_spec) ) result_df = df_with_target.filter(F.col("Time") >= F.col("target_time")).drop("target_time") # 展示结果 result_df.show()
原代码问题分析
你之前的实现错误在于窗口函数的用法:
# 错误写法:无聚合函数,仅返回当前行的Time,且仅LatestRecord=1的行有值 F.col('Time').over(Window.partitionBy('ID'))
这会导致earliestDate列只有符合条件的行有值,其余均为null,过滤时null >= Time的判断不成立,最终无法得到正确结果。
内容的提问来源于stack exchange,提问作者John Stud
相关产品推荐
相关产品推荐

