如何在Spark中基于起止时间戳关联两个DataFrame并生成连续行程
合并行程DataFrame生成连续往返记录
问题背景
现有两个记录公交行程的DataFrame:
- Source DataFrame:记录公交从CITY1出发的时间
- Destination DataFrame:记录公交抵达CITY2的时间
需要将这两个DataFrame合并,生成连续的行程链:每段行程的抵达时间作为下一段行程的出发时间,公交在CITY1和CITY2之间往返,最终得到包含完整往返记录的结果。
原始数据结构
Source DataFrame
+----+-------------------+------+ | ID|START_TIMESTAMP |SOURCE| +----+-------------------+------+ |BUS1|2023-12-17 07:27:00| CITY1| |BUS1|2023-12-17 09:50:00| CITY1| |BUS1|2023-12-17 12:23:00| CITY1| |BUS1|2023-12-17 16:13:00| CITY1| |BUS2|2023-12-17 06:08:00| CITY1| |BUS2|2023-12-17 08:31:00| CITY1| |BUS2|2023-12-17 11:13:00| CITY1| |BUS2|2023-12-17 14:44:00| CITY1| +----+-------------------+------+
Destination DataFrame
+----+-------------------+-----------+ | ID| END_TIMESTAMP |DESTINATION| +----+-------------------+-----------+ |BUS1|2023-12-17 08:27:00| CITY2 | |BUS1|2023-12-17 11:13:00| CITY2 | |BUS1|2023-12-17 14:50:00| CITY2 | |BUS2|2023-12-17 07:09:00| CITY2 | |BUS2|2023-12-17 09:53:00| CITY2 | |BUS2|2023-12-17 13:31:00| CITY2 | +----+-------------------+-----------+
期望输出
+----+-------------------+-----------------------+------+ | ID|START_TIMESTAMP | END_TIMESTAMP |stop | +----+-------------------+-----------------------+------+ |BUS1|2023-12-17 07:27:00| 2023-12-17 08:27:00|CITY2 | |BUS1|2023-12-17 08:27:00| 2023-12-17 09:50:00|CITY1 | |BUS1|2023-12-17 09:50:00| 2023-12-17 11:13:00|CITY2 | |BUS1|2023-12-17 11:13:00| 2023-12-17 12:23:00|CITY1 | |BUS1|2023-12-17 12:23:00| 2023-12-17 14:50:00|CITY2 | |BUS1|2023-12-17 14:50:00| 2023-12-17 16:13:00|CITY1 | |BUS2|2023-12-17 06:08:00| 2023-12-17 07:09:00|CITY2 | |BUS2|2023-12-17 07:09:00| 2023-12-17 08:31:00|CITY1 | |BUS2|2023-12-17 08:31:00| 2023-12-17 09:53:00|CITY2 | |BUS2|2023-12-17 09:53:00| 2023-12-17 11:13:00|CITY1 | |BUS2|2023-12-17 11:13:00| 2023-12-17 13:31:00|CITY2 | |BUS2|2023-12-17 13:31:00| 2023-12-17 14:44:00|CITY1 | +----+-------------------+-----------------------+------+
解决方案(PySpark实现)
步骤说明
- 对两个DataFrame按
ID分组,分别按时间字段升序排序,添加行号用于关联 - 关联正向行程(CITY1→CITY2):将Source的出发时间和Destination的到达时间匹配
- 生成返程行程(CITY2→CITY1):用上一段的到达时间作为出发时间,下一段的出发时间作为到达时间
- 合并正向和返程行程,按
ID和START_TIMESTAMP排序得到最终结果
代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("BusTripMerge").getOrCreate() # 构建Source DataFrame source_data = [ ("BUS1", "2023-12-17 07:27:00", "CITY1"), ("BUS1", "2023-12-17 09:50:00", "CITY1"), ("BUS1", "2023-12-17 12:23:00", "CITY1"), ("BUS1", "2023-12-17 16:13:00", "CITY1"), ("BUS2", "2023-12-17 06:08:00", "CITY1"), ("BUS2", "2023-12-17 08:31:00", "CITY1"), ("BUS2", "2023-12-17 11:13:00", "CITY1"), ("BUS2", "2023-12-17 14:44:00", "CITY1"), ] source_df = spark.createDataFrame(source_data, ["ID", "START_TIMESTAMP", "SOURCE"]) \ .withColumn("START_TIMESTAMP", F.to_timestamp("START_TIMESTAMP")) # 构建Destination DataFrame dest_data = [ ("BUS1", "2023-12-17 08:27:00", "CITY2"), ("BUS1", "2023-12-17 11:13:00", "CITY2"), ("BUS1", "2023-12-17 14:50:00", "CITY2"), ("BUS2", "2023-12-17 07:09:00", "CITY2"), ("BUS2", "2023-12-17 09:53:00", "CITY2"), ("BUS2", "2023-12-17 13:31:00", "CITY2"), ] dest_df = spark.createDataFrame(dest_data, ["ID", "END_TIMESTAMP", "DESTINATION"]) \ .withColumn("END_TIMESTAMP", F.to_timestamp("END_TIMESTAMP")) # 1. 给Source和Destination添加分组内的行号 source_window = Window.partitionBy("ID").orderBy("START_TIMESTAMP") source_ranked = source_df.withColumn("rank", F.row_number().over(source_window)) dest_window = Window.partitionBy("ID").orderBy("END_TIMESTAMP") dest_ranked = dest_df.withColumn("rank", F.row_number().over(dest_window)) # 2. 关联正向行程(CITY1→CITY2) forward_trips = source_ranked.join(dest_ranked, on=["ID", "rank"], how="inner") \ .select( "ID", "START_TIMESTAMP", "END_TIMESTAMP", F.col("DESTINATION").alias("stop") ) # 3. 生成返程行程(CITY2→CITY1) # 获取每个ID的下一段出发时间 source_with_next = source_ranked.withColumn( "NEXT_START", F.lead("START_TIMESTAMP").over(source_window) ).filter(F.col("NEXT_START").isNotNull()) # 关联得到返程行程 return_trips = dest_ranked.join(source_with_next, on=["ID", "rank"], how="inner") \ .select( "ID", F.col("END_TIMESTAMP").alias("START_TIMESTAMP"), F.col("NEXT_START").alias("END_TIMESTAMP"), F.lit("CITY1").alias("stop") ) # 4. 合并正向和返程行程,排序 final_trips = forward_trips.union(return_trips) \ .orderBy("ID", "START_TIMESTAMP") # 显示结果 final_trips.show(truncate=False)
结果验证
运行上述代码后,输出的DataFrame将与期望结构完全一致,每个公交的行程按时间顺序连续排列,往返于CITY1和CITY2之间。
内容的提问来源于stack exchange,提问作者RMK
相关产品推荐
相关产品推荐

