PySpark DataFrame过滤需求:移除每个ID首个非空value前的行
PySpark数据处理需求
现有DataFrame定义
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, FloatType, TimestampType from pyspark.sql.functions import col, min spark = SparkSession.builder.appName("example").getOrCreate() # 定义DataFrame schema schema = StructType([ StructField("dt", StringType(), True), StructField("id", StringType(), True), StructField("value", FloatType(), True) ]) # 创建DataFrame data = [ ('2023-04-24 00:00:00', 'A', None), ('2023-04-25 00:00:00', 'A', 100.5), ('2023-04-24 00:00:00', 'B', None), ('2023-04-25 00:00:00', 'B', None), ('2023-04-26 00:00:00', 'A', 110.0), ('2023-04-26 00:00:00', 'B', None), ('2023-04-27 00:00:00', 'A', None), ('2023-04-27 00:00:00', 'B', 50.5), ('2023-04-28 00:00:00', 'A', 105.0), ('2023-04-28 00:00:00', 'B', None), ('2023-04-29 00:00:00', 'B', 55.5), ('2023-04-29 00:00:00', 'A', 107.0) ] df = spark.createDataFrame(data, schema) df = df.withColumn("dt", col("dt").cast(TimestampType()))
处理需求
按id分组,移除每个分组中首个value非空行之前的所有记录,保留该行及之后的所有行(包括后续value为空的行)。
预期结果
+-------------------+---+-----+ | dt| id|value| +-------------------+---+-----+ |2023-04-25 00:00:00| A|100.5| |2023-04-26 00:00:00| A|110.0| |2023-04-27 00:00:00| A| null| |2023-04-28 00:00:00| A|105.0| |2023-04-29 00:00:00| A|107.0| |2023-04-27 00:00:00| B| 50.5| |2023-04-28 00:00:00| B| null| |2023-04-29 00:00:00| B| 55.5| +-------------------+---+-----+
解决方案
通过分组计算每个id的首个非空value对应时间,再过滤掉早于该时间的行即可实现:
# 1. 计算每个id首次出现非空value的时间 first_valid_dt = df.filter(col("value").isNotNull()) \ .groupBy("id") \ .agg(min("dt").alias("first_valid_dt")) # 2. 关联原表并过滤数据 result_df = df.join(first_valid_dt, on="id", how="inner") \ .filter(col("dt") >= col("first_valid_dt")) \ .drop("first_valid_dt") # 查看最终结果 result_df.orderBy("id", "dt").show()
逻辑说明
- 第一步:先筛选所有
value非空的行,按id分组后取每组最早的dt,即该id首次出现有效value的时间点。 - 第二步:将原表与第一步的结果关联,只保留
dt大于等于首次有效时间的行,最后删除辅助列得到目标数据。
内容的提问来源于stack exchange,提问作者andKaae
相关产品推荐
相关产品推荐

