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

如何在PySpark中按源和目标分组计算跨多行的时间差

需求实现方案

需求说明

计算同一ID下首个leg_sent时间与最后一个leg_received时间的差值,并按Source和Destination分组统计相关指标。数据集test_subset中,ID为传输标识,每行对应一段传输路径(leg),同一ID的Source和Destination固定,传输段数量为2-10个。

示例数据

IDFromToleg_sentleg_receivedSourceDestination
1btrABCXYZ08:22:2308:22:41GBFR
1btrXYZDEF08:22:4908:23:05GBFR
2vyuLMNJFK14:35:1114:35:23USDE
2vyuJFKHIJ14:35:3514:35:48USDE
2vyuHIJTPQ14:35:5114:36:25USDE

修改后代码

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"))  # 耗时范围(秒)
)

代码说明

  1. ID级聚合:先按ID分组,提取每个完整传输的首个发送时间(最早的leg_sent)和最后接收时间(最晚的leg_received),同时获取固定的Source和Destination。
  2. 总耗时计算:将时间字符串转换为Unix时间戳(秒)后相减,得到单次传输的总耗时。若你的时间字段是timestamp类型,建议使用F.timestamp_diff("last_leg_received", "first_leg_sent", "seconds"),避免时间格式解析问题。
  3. 路径分组统计:基于ID级的总耗时数据,按Source和Destination分组,统计传输次数、平均/最小/最大总耗时,以及耗时范围。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 04:30:41