PySpark按时间区间关联DataFrame的实现方法
Hey, let's work through this problem where we need to join two DataFrames on the Fruit field, and only match records where the actual consumption time is within 5 minutes of the suggested time—plus we have to keep all rows from df_fruit (left join).
First, the key challenge here is comparing time strings directly isn't straightforward. So we'll convert those "HH:mm" strings into total minutes since midnight, which makes calculating time differences a breeze.
Step 1: Convert time strings to total minutes
We'll turn times like "10:00" into 600 (1060) and "12:35" into 755 (1260+35). This lets us easily check if the absolute difference between suggested and actual time is ≤5 minutes.
You can use either a custom UDF or Spark's built-in functions for this. I'll show both approaches below.
Step 2: Perform the left join with time condition
We'll use Spark's join method, specify how="left" to retain all df_fruit rows, and add the dual condition: matching Fruit values, and time difference ≤5 minutes.
Full Code Implementation
Option 1: Using a custom UDF (intuitive for beginners)
from pyspark.sql import SparkSession from pyspark.sql.functions import col, abs, udf from pyspark.sql.types import IntegerType # Initialize SparkSession spark = SparkSession.builder.appName("FruitTimeMatch").getOrCreate() # Create original DataFrames df_fruit = spark.createDataFrame( [("Apple", "10:00"),("Orange", "12:35"),("Apple", "11:36"),("Apple","12:48"),("Pear","11:00")], ["Fruit", "Time"] ) df_calories = spark.createDataFrame( [("Apple", "10:02", "86g", "1cal"),("Orange", "12:39", "75g", "14cal"),("Apple", "10:04", "9g", "47cal"),("Apple","12:46", "25g", "9cal"),("Orange","12:33", "75g", "2cal")], ["Fruit", "Time", "Weight", "Calories"] ) # UDF to convert "HH:mm" to total minutes def time_to_total_minutes(time_str): hours, mins = map(int, time_str.split(":")) return hours * 60 + mins time_to_min_udf = udf(time_to_total_minutes, IntegerType()) # Add minute columns to both DataFrames df_fruit = df_fruit.withColumn("suggested_min", time_to_min_udf(col("Time"))) df_calories = df_calories.withColumn("actual_min", time_to_min_udf(col("Time"))) # Left join with time condition joined_df = df_fruit.join( df_calories, (df_fruit["Fruit"] == df_calories["Fruit"]) & (abs(col("suggested_min") - col("actual_min")) <= 5), how="left" ) # Select and rename columns to match expected output result = joined_df.select( df_fruit["Fruit"], df_fruit["Time"].alias("Suggested_Time"), df_calories["Time"].alias("Actual_Time"), df_calories["Weight"], df_calories["Calories"] ) # Show the result result.show()
Option 2: Using Spark built-in functions (no UDF needed)
If you prefer to avoid UDFs (they can be less efficient at scale), use Spark's timestamp functions:
from pyspark.sql.functions import to_timestamp, hour, minute # Convert time to timestamp and calculate total minutes df_fruit = df_fruit.withColumn( "suggested_min", hour(to_timestamp(col("Time"), "HH:mm")) * 60 + minute(to_timestamp(col("Time"), "HH:mm")) ) df_calories = df_calories.withColumn( "actual_min", hour(to_timestamp(col("Time"), "HH:mm")) * 60 + minute(to_timestamp(col("Time"), "HH:mm")) ) # Same join and selection as Option 1 joined_df = df_fruit.join( df_calories, (df_fruit["Fruit"] == df_calories["Fruit"]) & (abs(col("suggested_min") - col("actual_min")) <= 5), how="left" ) result = joined_df.select( df_fruit["Fruit"], df_fruit["Time"].alias("Suggested_Time"), df_calories["Time"].alias("Actual_Time"), df_calories["Weight"], df_calories["Calories"] ) result.show()
Expected Output
Both options will produce the output you're looking for:
+------+-------------+-----------+------+--------+ | Fruit|Suggested_Time|Actual_Time|Weight|Calories| +------+-------------+-----------+------+--------+ | Apple| 10:00| 10:02| 86g| 1cal| | Apple| 10:00| 10:04| 9g| 47cal| |Orange| 12:35| 12:39| 75g| 14cal| |Orange| 12:35| 12:33| 75g| 2cal| | Apple| 11:36| null| null| null| | Apple| 12:48| 12:46| 25g| 9cal| | Pear| 11:00| null| null| null| +------+-------------+-----------+------+--------+
Key Notes
- The left join ensures every row from
df_fruitis kept, even if there's no matching consumption record (like Apple 11:36 and Pear 11:00). - Using built-in functions is better for performance in large datasets, as UDFs can't be optimized by Spark's Catalyst optimizer as well.
内容的提问来源于stack exchange,提问作者Bryce Ramgovind

