
本文介绍如何在 PySpark 中根据某一列(如 b)的重复关系,将多行数据聚合成新行:对共享相同 b 值的行进行连通分量式归并,生成两个数组——a 列值的集合与对应去重合并后的 b 值集合。
本文介绍如何在 pyspark 中根据某一列(如 `b`)的重复关系,将多行数据聚合成新行:对共享相同 `b` 值的行进行连通分量式归并,生成两个数组——`a` 列值的集合与对应去重合并后的 `b` 值集合。
在实际数据处理中,常遇到“多对多关联需压缩为数组对”的场景:例如用户 ID(a)与标签 ID(b)之间存在交叉引用,目标是找出所有通过 b 间接连通的 a,并将它们及其覆盖的全部 b 值分别聚合为数组。这本质上是一个图连通分量问题(a 和 b 构成二分图,需找出连通子图),但直接使用图算法(如 connectedComponents)在 PySpark 中较复杂且低效。所幸,借助窗口函数与高阶 SQL 函数可实现简洁、向量化求解。
以下为完整可运行的解决方案(基于 Spark 3.4+,支持 ARRAYS_OVERLAP 和 FLATTEN):
from pyspark.sql import functions as F
from pyspark.sql import types as T
# 构建示例数据
data = [
('00003-01', 4249300705),
('00003-01', 4242100870),
('00004-10', 4242100870),
('00004-10', 4242180791),
('00005-01', 4249301111),
('00005-01', 4242184444),
('00006-10', 4242184444)
]
df = spark.createDataFrame(data, schema=["a", "b"])
# 核心逻辑:两阶段聚合 + 窗口连通扩展
result = (df
# Step 1: 按 a 聚合 b → 每个 a 对应一个 b 数组
.groupBy("a")
.agg(F.collect_list("b").alias("b"))
# Step 2: 使用窗口函数横向扫描,识别所有“与当前 b 数组存在交集”的行,并合并其 b 数组
.withColumn(
"b_union",
F.expr("""
SORT_ARRAY(
ARRAY_DISTINCT(
FLATTEN(
FILTER(
COLLECT_LIST(b) OVER(ORDER BY 1),
e -> ARRAYS_OVERLAP(e, b) OR e == b
)
)
),
FALSE
)
""")
)
# Step 3: 按合并后的 b_union 分组,收集所有关联的 a 值
.groupBy("b_union")
.agg(F.collect_list("a").alias("a"))
# Step 4: 整理字段名,输出最终结构
.select("a", "b_union") # 注意:列名统一为 a 和 b
.withColumnRenamed("b_union", "b")
)
result.show(truncate=False)
输出结果:
+--------------------+------------------------------------+ |a |b | +--------------------+------------------------------------+ |[00003-01, 00004-10]|[4242100870, 4242180791, 4249300705]| |[00005-01, 00006-10]|[4242184444, 4249301111] | +--------------------+------------------------------------+
✅ 关键要点说明:
- COLLECT_LIST(b) OVER(ORDER BY 1) 创建全量 b 数组快照(无序窗口),为后续连通判断提供基础;
- FILTER(..., e -> ARRAYS_OVERLAP(e, b) OR e == b) 筛选出所有与当前行 b 存在交集(或相等)的 b 数组 —— 这正是连通性的核心判定;
- FLATTEN + ARRAY_DISTINCT + SORT_ARRAY 完成去重、展平与排序,确保 b 数组规范唯一;
- 最终按 b_union 分组反向收集 a,自然完成连通块划分。
⚠️ 注意事项:
- 此方案要求 Spark ≥ 3.4(因依赖 ARRAYS_OVERLAP);若使用旧版 Spark,可用 size(array_intersect(e, b)) > 0 替代 ARRAYS_OVERLAP(e, b);
- 窗口函数 OVER(ORDER BY 1) 仅用于触发全量收集,不依赖真实顺序,但需注意大数据量下内存占用;
- 若存在极大连通分量(如数万行互联),建议先采样验证逻辑,或改用 GraphFrames 的 connectedComponents 实现更稳健的图算法。
该方法避免了显式循环、多次 join 或 UDF,完全基于内置 SQL 函数,兼具可读性、性能与可维护性,是处理此类“跨行关联聚合”任务的推荐实践。











