PySpark DataFrame按ID与日期规则筛选行的实现求助
解决方案
要实现每个id下最早日期保留id_index=1、下一个日期保留id_index=2的需求,可通过窗口函数生成日期序号再匹配筛选完成,具体步骤如下:
步骤1:为每个ID的日期生成顺序序号
用row_number()窗口函数,按id分组、date升序排序,给每个id下的日期分配递增序号(最早日期为1,后续依次+1)。
步骤2:匹配序号与id_index筛选行
直接过滤出id_index等于日期序号的行,即可得到目标结果。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number # 初始化SparkSession spark = SparkSession.builder.appName("id_date_index_match").getOrCreate() # 模拟输入数据 data = [ (1, "2023-05-14", 107.27966, 1), (1, "2023-05-14", 78.23225, 2), (1, "2023-05-14", 226.91467, 3), (1, "2023-05-15", 107.27966, 1), (1, "2023-05-15", 78.23225, 2), (1, "2023-05-15", 226.91467, 3), (2, "2023-05-14", 249.05295, 1), (2, "2023-05-14", 2.0, 2), (2, "2023-05-14", 0.0, 3), (2, "2023-05-15", 249.05295, 1), (2, "2023-05-15", 2.0, 2), (2, "2023-05-15", 0.0, 3) ] df = spark.createDataFrame(data, ["id", "date", "value", "id_index"]) # 定义窗口:按id分组,date升序排序 window_spec = Window.partitionBy("id").orderBy("date") # 添加日期序号列 df_with_seq = df.withColumn("date_seq", row_number().over(window_spec)) # 筛选匹配的行 result_df = df_with_seq.filter(df_with_seq.id_index == df_with_seq.date_seq).drop("date_seq") # 展示结果 result_df.show()
执行结果
+---+----------+---------+--------+ | id| date| value|id_index| +---+----------+---------+--------+ | 1|2023-05-14|107.27966| 1| | 1|2023-05-15| 78.23225| 2| | 2|2023-05-14|249.05295| 1| | 2|2023-05-15| 2.0| 2| +---+----------+---------+--------+
该方案核心是通过窗口函数将日期顺序转化为可匹配的数值,直接与id_index做等值筛选,逻辑清晰且性能高效,避免了复杂的索引创建和去重操作。
内容的提问来源于stack exchange,提问作者orocd398897
相关产品推荐
相关产品推荐

