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

如何使用Pyspark窗口函数统计30分钟分组下的站点间换乘次数

问题:Pyspark实现连续站点换乘组合统计

我正在使用Pyspark,需要开发实现如下功能的函数:

现有描述火车用户出行交易的数据集,数据结构如下:

+----+-------------------+-------+---------------+-------------+-------------+
|USER|       DATE        |LINE_ID|      STOP     | TOPOLOGY_ID |TRANSPORT_ID |
+------------------------+-------+---------------+-------------+-------------+
|John|2021-01-27 07:27:34|      7| King Cross    |       171235|       03    |
|John|2021-01-27 07:28:00|     40| White Chapell |       123582|       03    |  
|John|2021-01-27 07:35:30|      4| Reaven        |       171565|       03    |  
|Tom |2021-01-27 07:27:23|      7| King Cross    |       171235|       03    |    
|Tom |2021-01-27 07:28:30|     40| White Chapell |       123582|       03    |                   
+----+-------------------+-------+---------------+-------------+-------------+

我需要统计30分钟时间分组内,A-B、B-C等连续站点换乘组合的出现次数。举例说明:用户John于7:27从King Cross站前往White Chapell站,7:35再前往Reaven站;同时用户Tom于7:27从King Cross站前往White Chapell站,7:32前往Oxford Circus站。

最终需要输出的结果结构如下:

+----------------------+-----------------+---------------+-----------+
|          DATE        |   ORIG_STOP     |   DEST_STOP   | NUM_TRANS |
+----------------------+-----------------+---------------+-----------+
|   2021-01-27 07:00:00|  King Cross     | White Chapell |       2   |
|   2021-01-27 07:30:00|  White Chapell  | Reaven        |       1   |              
+----------------------+-----------------+---------------+-----------+

我已经尝试使用窗口函数实现,但未得到预期结果,请问该如何实现该需求?


实现方案

核心逻辑是先为每个用户的出行记录匹配相邻行程,得到起点-终点换乘对,再按30分钟时间窗口分组统计次数,完整实现代码如下:

步骤1:导入依赖

from pyspark.sql.functions import col, lead, count
from pyspark.sql.window import Window

步骤2:匹配相邻行程构造换乘对

按用户分区、交易时间排序,用lead窗口函数取当前站的下一站、下一站交易时间:

# 定义用户行程排序窗口
user_trip_window = Window.partitionBy("USER").orderBy("DATE")

# 构造换乘对,过滤无后续行程的最后一条记录
trip_pair_df = df.withColumn("DEST_STOP", lead("STOP", 1).over(user_trip_window)) \
                 .withColumn("ARRIVAL_DATE", lead("DATE", 1).over(user_trip_window)) \
                 .filter(col("DEST_STOP").isNotNull())

步骤3:生成30分钟时间分组

将时间对齐到最近的30分钟整点(7:00-7:29归为7:00分组,7:30-7:59归为7:30分组,和示例规则对齐采用到达时间计算窗口):

# 1800秒=30分钟,计算所属30分钟窗口
trip_pair_df = trip_pair_df.withColumn(
    "DATE",
    col("ARRIVAL_DATE") - (col("ARRIVAL_DATE").cast("long") % 1800).cast("interval seconds")
)

若业务规则需要按出发时间计算窗口,将上面代码中的ARRIVAL_DATE替换为原字段DATE即可。

步骤4:分组统计换乘次数

result = trip_pair_df.groupBy(
    "DATE",
    col("STOP").alias("ORIG_STOP"),
    "DEST_STOP"
).agg(
    count("*").alias("NUM_TRANS")
).orderBy("DATE", "ORIG_STOP")

补充说明

如果需要区分不同线路/车次的行程,可以将LINE_ID、TRANSPORT_ID加入partitionBy的分区条件,避免把不同车次的行程误判为连续换乘。


内容的提问来源于stack exchange,提问作者tomruarol

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 07:39:00