如何在PySpark中按源和目标分组计算跨多行的时间差
需求实现方案
需求说明
计算同一ID下首个leg_sent时间与最后一个leg_received时间的差值,并按Source和Destination分组统计相关指标。数据集test_subset中,ID为传输标识,每行对应一段传输路径(leg),同一ID的Source和Destination固定,传输段数量为2-10个。
示例数据
| ID | From | To | leg_sent | leg_received | Source | Destination |
|---|---|---|---|---|---|---|
| 1btr | ABC | XYZ | 08:22:23 | 08:22:41 | GB | FR |
| 1btr | XYZ | DEF | 08:22:49 | 08:23:05 | GB | FR |
| 2vyu | LMN | JFK | 14:35:11 | 14:35:23 | US | DE |
| 2vyu | JFK | HIJ | 14:35:35 | 14:35:48 | US | DE |
| 2vyu | HIJ | TPQ | 14:35:51 | 14:36:25 | US | DE |
修改后代码
from pyspark.sql import functions as F # 1. 按ID聚合,获取单次完整传输的关键时间点与路径信息 id_transfer_stats = ( test_subset .groupBy("ID") .agg( F.min("leg_sent").alias("first_leg_sent"), # 同一ID的首个发送时间 F.max("leg_received").alias("last_leg_received"), # 同一ID的最后接收时间 F.first("Source").alias("Source"), # 同一ID的Source固定,取任意值即可 F.first("Destination").alias("Destination") # 同一ID的Destination固定,取任意值即可 ) # 计算单次传输的总耗时(秒):若时间为timestamp类型,可改用F.timestamp_diff .withColumn("total_transfer_time", F.unix_timestamp("last_leg_received", "HH:mm:ss") - F.unix_timestamp("first_leg_sent", "HH:mm:ss")) ) # 2. 按Source和Destination分组,统计整体传输指标 final_grouped_stats = ( id_transfer_stats .groupBy("Source", "Destination") .agg( F.count("ID").alias("num_transfers"), # 该路径对的传输总次数 F.avg("total_transfer_time").alias("avg_total_time"), # 平均总耗时(秒) F.min("total_transfer_time").alias("min_total_time"), # 最小总耗时(秒) F.max("total_transfer_time").alias("max_total_time") # 最大总耗时(秒) ) .withColumn("time_range", F.col("max_total_time") - F.col("min_total_time")) # 耗时范围(秒) )
代码说明
- ID级聚合:先按ID分组,提取每个完整传输的首个发送时间(最早的leg_sent)和最后接收时间(最晚的leg_received),同时获取固定的Source和Destination。
- 总耗时计算:将时间字符串转换为Unix时间戳(秒)后相减,得到单次传输的总耗时。若你的时间字段是
timestamp类型,建议使用F.timestamp_diff("last_leg_received", "first_leg_sent", "seconds"),避免时间格式解析问题。 - 路径分组统计:基于ID级的总耗时数据,按Source和Destination分组,统计传输次数、平均/最小/最大总耗时,以及耗时范围。
内容的提问来源于stack exchange,提问作者Qaribbean
相关产品推荐
相关产品推荐

