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
相关产品推荐
相关产品推荐

