基于日期区间的Left Join:为DataFrame添加tag列的实现问题
问题解决:为DataFrame添加符合日期区间规则的tag列
你的代码存在以下问题:
- 列名拼写错误:将
Date 1写成了Date1,exit_date写成了exist_date - 逻辑运算优先级问题:
&和|的优先级高于比较运算符,未加括号会导致逻辑判断出错 - 多匹配未处理:df2中同一user_id存在多条记录,直接join会导致df1的原始记录被重复,需要先聚合判断是否存在满足条件的区间
正确解决方案
以下是两种可行的实现方式,均能得到你预期的结果:
方法一:先聚合df2的匹配结果再关联df1
from pyspark.sql import functions as F # 关联df1和df2,标记每条区间记录是否匹配当前user_id的Date 1 df2_match = df2.join(df1, on="user_id", how="inner") \ .withColumn( "is_match", F.when( (F.col("Date 1") >= F.col("entrance_date")) & (F.col("exit_date").isNull() | (F.col("Date 1") <= F.col("exit_date"))), F.lit(True) ).otherwise(F.lit(False)) ) # 按user_id聚合,只要有一条匹配就标记为True user_tag = df2_match.groupBy("user_id") \ .agg(F.max("is_match").alias("tag")) # 关联回df1,补全无匹配的user_id为False result_df = df1.join(user_tag, on="user_id", how="left") \ .fillna(False, subset=["tag"])
方法二:使用窗口函数处理多匹配
from pyspark.sql import functions as F from pyspark.sql.window import Window # 左连接df1和df2 joined_df = df1.join(df2, on="user_id", how="left") # 添加单条区间的匹配标记 joined_df = joined_df.withColumn( "is_match", F.when( (F.col("Date 1") >= F.col("entrance_date")) & (F.col("exit_date").isNull() | (F.col("Date 1") <= F.col("exit_date"))), F.lit(True) ).otherwise(F.lit(False)) ) # 按user_id取全局匹配结果,去重得到最终df1 window = Window.partitionBy("user_id") result_df = joined_df.withColumn("tag", F.max("is_match").over(window)) \ .select("user_id", "Date 1", "tag") \ .distinct()
最终结果
执行上述代码后,result_df会输出:
| user_id | Date 1 | tag |
|---|---|---|
| 1 | 2023-01-01 | True |
| 2 | 2020-02-15 | False |
| 3 | 2022-03-02 | True |
内容的提问来源于stack exchange,提问作者f.ivy
相关产品推荐
相关产品推荐

