基于自定义函数[next Date]实现PySpark DataFrame日期匹配INNER JOIN
Got it, let's tackle this problem step by step. You need to perform an inner join between DataFrames a and b where we match on user and id, and also find the latest (maximum) A_Date from a that is strictly less than B_Date from b for each matching pair. Here's how to do it efficiently with PySpark window functions:
Step 1: Import Required Functions
First, bring in the PySpark utilities we'll need:
from pyspark.sql import Window from pyspark.sql.functions import col, row_number
Step 2: Initial Join & Filter
Start by joining the two DataFrames on user and id, then filter to keep only rows where A_Date falls before B_Date:
# Inner join on matching user and id joined_df = a.join(b, on=["user", "id"], how="inner") # Keep only valid date pairs where A_Date is earlier than B_Date filtered_df = joined_df.filter(col("A_Date") < col("B_Date"))
Step 3: Use Window Function to Get Closest Matching Date
We'll use a window partitioned by user, id, and B_Date (to group each row from b with its potential matches from a). Then we sort A_Date in descending order and pick the top row—this gives us the closest date in a that's before B_Date:
# Define window to group by each unique (user, id, B_Date) and sort A_Date descending window_spec = Window.partitionBy(col("user"), col("id"), col("B_Date")) \ .orderBy(col("A_Date").desc()) # Rank rows in each partition and keep only the top (latest) matching date ranked_df = filtered_df.withColumn("row_rank", row_number().over(window_spec)) \ .filter(col("row_rank") == 1) \ .drop("row_rank")
Step 4: Reshape to Target DataFrame
Finally, select and reorder columns to match your desired output c:
c = ranked_df.select(col("B_Date"), col("user"), col("id"), col("info"), col("A_Date")) # Check the result c.show()
Sample Output
Running c.show() will return exactly the DataFrame you specified:
+----------+------+---+----+----------+ | B_Date | user | id|info| A_Date | +----------+------+---+----+----------+ |2017-05-15| foo| 24| def|2017-05-14| |2017-05-22| user2| 27| mno|2017-05-19| |2017-05-20| user3| 28| mno|2017-05-19| +----------+------+---+----+----------+
Why This Works
- The initial join ensures we only consider rows where
userandidmatch betweenaandb. - The filter narrows down to valid date pairs where
a's date is earlier thanb's. - The window function lets us isolate each unique entry from
b, sort its matching dates fromaby recency, and pick the closest valid date—this avoids redundant matches and gives exactly the result you need.
内容的提问来源于stack exchange,提问作者Shekhar Koirala

