You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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实现)

步骤说明

  1. 对两个DataFrame按ID分组,分别按时间字段升序排序,添加行号用于关联
  2. 关联正向行程(CITY1→CITY2):将Source的出发时间和Destination的到达时间匹配
  3. 生成返程行程(CITY2→CITY1):用上一段的到达时间作为出发时间,下一段的出发时间作为到达时间
  4. 合并正向和返程行程,按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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 10:17:33