如何使用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
相关产品推荐
相关产品推荐

