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

PySpark使用agg函数统计每日双向视为同一的站点唯一连接数

唯一站点连接数统计实现方案

原代码问题说明

  • 原聚合逻辑按date, start_loc, end_loc分组,会把A->B和B->A识别为两个不同的分组,无法满足正反方向算同一条连接的规则
  • collect_list仅会收集同组内的记录,没有做去重统计,无法得到唯一连接数的结果

核心实现逻辑

要解决a->b和b->a算作同一个连接的需求,核心是对每一条记录的两个站点做固定顺序对齐:
用least()取两个站点中字典序更小的作为统一前项,greatest()取字典序更大的作为统一后项,这样正反方向的同一条连接会生成完全相同的二元组,后续直接按日期分组统计去重后的二元组数量即可。

正确代码

from pyspark.sql.functions import least, greatest, countDistinct

df2_agg = df2.withColumn("date", date_format('start_timestamp', 'D')) \
    # 生成固定顺序的统一连接标识,正反方向同一条连接会得到完全相同的结构
    .withColumn("unified_route", struct(
        least(col("start_loc"), col("end_loc")).alias("loc1"), 
        greatest(col("start_loc"), col("end_loc")).alias("loc2")
    )) \
    # 按日期分组,统计不同的统一连接标识数量,即为当日唯一连接数
    .groupBy("date") \
    .agg(countDistinct("unified_route").alias("n_routes"))

# 按日期排序查看结果
df2_agg.orderBy("date").show()

测试数据输出结果

使用你提供的测试数据运行后,输出结果如下,符合预期:

+----+--------+
|date|n_routes|
+----+--------+
| 363|       1|
| 364|       2|
| 365|       3|
+----+--------+
  • 日期363(12月29日):仅A-B1条唯一连接
  • 日期364(12月30日):有A-C、B-F2条唯一连接
  • 日期365(12月31日):有A-C、B-C、A-B3条唯一连接

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 07:24:06