Connected Component from DFS to BigQuery

连通分量计算的算法演进之路

去年迁移地图处理 pipeline 的时候,遇到了一段道路连通分量计算的代码。原本这部分的处理用的是 GraphFrame 的图计算库

GraphFrames 是一个用于 Apache Spark 的软件包,它提供基于 DataFrame 的图功能。它提供 Scala、Java 和 Python 的高级 API: https://graphframes.io/

,非常简单好用。但迁移后的 BigQuery 的官方文档中并没有图相关的函数,就随便让 Gemini 帮我糊了一个 SQL 实现,跑了一下感觉还行就直接用了。

原本的代码大概是这样的:

import org.apache.spark.sql.functions._
import org.graphframes.GraphFrame

// 准备 node 和 edge
val roads = spark.table("road").select("src", "dst", "oid")
val edges = roads.union(roads.select("dst", "src", "oid").toDF("src", "dst", "oid"))
val nodes = edges.select("src").distinct().union(edges.select("dst").distinct()).toDF("id")

// 计算连通分量
val g = GraphFrame(nodes, edges)
val result = g.connectedComponents.setUseLocalCheckpoints(true).run()

// 将结果与原始道路表连接,得到每条道路所属的连通分量
val road_x_connected_component = result.join(roads, expr("id=dst"))
  .union(result.join(roads, expr("id=src")))
  .select("oid", "component")
  .dropDuplicates()

说实话这个代码真正计算连通分量的地方就只有 GraphFrame 封装好的一个函数,里面究竟干了什么我也不知道。Gemini 给我生成的 SQL 代码我也没太看明白,所以这里我再重新学习一下这个基础图问题的相关算法。

那么首先什么是连通分量呢?

给定一个无向图 G=(V,E)G = (V, E),其中 VV 是节点集合,EE 是边集合。一个连通分量是图中一个极大的节点子集,使得子集中的任意两个节点之间都存在一条路径,且无法再加入任何外部节点而保持这一性质。

换句话说,连通分量就是图中彼此”连通”的节点群组。每个节点恰好属于一个连通分量,所有连通分量构成了对整个图的一个划分。

一个直观的例子:假设有 6 个节点和以下 4 条边:

一个简单的无向图 一个简单的无向图
1 -- 2
2 -- 3
3 -- 4
5 -- 6

这个图绘制出来就会如右图所示,一共有两个连通分量:4 和 6。节点 1 可以经由 2、3 到达 4,但无论如何都无法到达 5 或 6。

所以连通分量计算的目标就是为每个节点分配一个标签(通常是该分量中最小的节点 ID),使得同一分量中的节点拥有相同的标签。

经典单机算法:DFS/BFS 与并查集

DFS/BFS 遍历

首先肯定不能绕过的是深度优先搜索(DFS)或广度优先搜索(BFS),感觉一下子梦回大学的数据结构课程。虽然现在上班好多年,已经不能手写代码了,但算法的核心思想还是记得的。这里稍微整理一下思路:

从一个未访问的节点出发,通过 DFS 或 BFS 遍历所有可达节点,将它们标记为同一分量。重复这个过程直到所有节点都被访问。

整体时间复杂度为 O(V+E)O(V + E),其中 VV 是节点数,EE 是边数。对于使用邻接表存储的图,这已经是最优的了。如果图是非常稠密的(EV2E \approx V^2),那么 O(V+E)O(V + E) 实际上就等同于 O(V2)O(V^2)

整体上在单机内存能容纳整个图的情况下非常高效,然而它本质上是一个顺序算法。整个遍历过程依赖于前一步的结果,难以并行化。如果图有数十亿个节点和边时,那单机的内存和计算能力就都扛不住了。

并查集(Union-Find / Disjoint Set Union)

Tarjan, R. E. (1975). “Efficiency of a Good But Not Linear Set Union Algorithm.” Journal of the ACM, 22(2), 215–225。

另一个经典方法是基于并查集数据结构。初始时每个节点自成一个集合。然后遍历每条边 (u,v)(u, v),将 uuvv 所在的集合合并(Union)。处理完所有边后,属于同一集合的节点就属于同一连通分量。

并查集支持两种核心操作:Find(找到节点所属集合的代表元素)和 Union(合并两个集合)。 通过路径压缩(Path Compression)和按秩合并(Union by Rank)两种优化,每次操作的均摊时间复杂度为 O(α(n))O(\alpha(n)),其中 α\alpha 是反阿克曼函数。

A(m,n)={n+1if m=0A(m1,1)if m>0 and n=0A(m1,A(m,n1))if m>0 and n>0{\displaystyle A(m,n)={\begin{cases} n+1 &\text{if } m=0\\ A(m-1,1) &\text{if } m>0 \text{ and } n=0\\ A(m-1,A(m,n-1)) &\text{if } m>0 \text{ and } n>0\\ \end{cases}} }

一个增长极其缓慢的函数,对所有实际输入规模都可以视为常数。

并查集的理论最优性分析参见 Fredman, M. & Saks, M. (1989). “The Cell Probe Complexity of Dynamic Data Structures.” STOC ‘89

并查集的优势在于它可以按任意顺序处理边,甚至支持流式(streaming)处理。但和 DFS/BFS 一样,它本质上仍是单机顺序算法,涉及大量的指针追踪(pointer chasing),这在分布式环境中效率极低——因为链表中的元素可能分布在不同的机器上。


三、并行计算的曙光:PRAM 算法(1980s–1990s)

随着并行计算理论的发展,研究者开始在 PRAM(Parallel Random Access Machine) 模型上探索连通分量的并行算法。PRAM 是一种理想化的并行计算模型,假设多个处理器可以同时访问一块共享内存。

3.1 Shiloach-Vishkin 算法(1982)

1982 年,Yossi Shiloach 和 Uzi Vishkin 提出了一个具有开创性意义的并行连通分量算法,发表在 Journal of Algorithms 上。这个算法在 CRCW(Concurrent Read Concurrent Write)PRAM 模型上运行,时间复杂度为 O(log n),使用 O(n + m) 个处理器。

Shiloach, Y. & Vishkin, U. (1982). “An O(log n) Parallel Connectivity Algorithm.” Journal of Algorithms, 3(1), 57–67。

Shiloach-Vishkin 算法的核心思想是将连通分量维护为一个有根树的森林,然后反复执行两种操作:

  1. Hooking(钩连): 通过图中的边,将一棵树”钩”到另一棵树上。具体来说,对于每条边 (u, v),如果 u 和 v 属于不同的树,就将其中一棵树的根指向另一棵树。为了避免形成环,算法精心设计了钩连的规则——只有在满足特定条件时才允许钩连。
  2. Shortcutting(快捷跳转)/ Pointer Jumping(指针跳跃): 将每个节点的父指针直接指向其祖父节点,从而快速压缩树的高度。经过 O(log n) 轮指针跳跃,所有节点都会直接指向各自的根节点。

这两个操作交替执行:Hooking 合并不同的树,Shortcutting 压缩树的高度,使得后续的 Hooking 更加高效。经过 O(log n) 轮迭代,所有属于同一连通分量的节点最终会汇聚到同一棵树中。

Shiloach-Vishkin 算法是后续几乎所有并行连通分量算法的思想源头。它引入的 Hooking + Shortcutting 范式成为了这个领域的基石。

3.2 后续的 PRAM 改进

Gazit, H. (1991). “An Optimal Randomized Parallel Algorithm for Finding Connected Components in a Graph.” SIAM J. Comput., 20(6), 1046–1067。

Karger, D. R., Nisan, N. & Parnas, M. (1992). “Fast Connected Components Algorithms for the EREW PRAM.” SPAA ‘92, 373–381。

Shiloach-Vishkin 之后,多位研究者在 PRAM 模型上提出了改进算法:

  • Awerbuch & Shiloach (1987) 简化了 Hooking 的判定规则,使算法更容易实现。
  • Gazit (1991) 提出了一个最优的随机化 PRAM 算法,在 CRCW PRAM 上达到 O(log n) 时间和 O((n + m) / log n) 的处理器数量。
  • Johnson & Metaxas (1991) 在更严格的 CREW(Concurrent Read Exclusive Write)PRAM 模型上给出了 O(log^1.5 n) 时间的确定性算法。
  • Karger, Nisan & Parnas (1992/1999) 在更严格的 EREW(Exclusive Read Exclusive Write)PRAM 模型上给出了随机化 O(log n) 时间的算法,以及确定性 O(log^1.5 n) 时间的算法。

然而,正如 Eppstein 和 Galil 所指出的,PRAM 模型”常被理论计算机科学家使用,却较少被实际并行机器的构建者采用”。PRAM 假设的共享内存模型与现实中分布式集群的架构有很大差距,这些算法很难直接迁移到实际系统中。


四、大数据时代:MapReduce 上的连通分量(2009–2014)

2004 年,Google 发表了具有划时代意义的 MapReduce 论文,开启了大数据分布式计算的新时代。MapReduce 框架将计算分为 Map 和 Reduce 两个阶段,天然适合处理大规模数据。然而,图算法(尤其是连通分量计算)在 MapReduce 上的实现面临独特的挑战:每轮 MapReduce 作业都有显著的启动和同步开销,因此减少迭代轮数成为最关键的优化目标。

4.1 Pegasus / Hash-Min(2009)

Kang, U., Tsourakakis, C. E. & Faloutsos, C. (2009). “PEGASUS: A Peta-Scale Graph Mining System — Implementation and Observations.” Ninth IEEE International Conference on Data Mining (ICDM)

2009 年,Kang、Tsourakakis 和 Faloutsos 在 IEEE ICDM 上发表了 Pegasus 系统,这是第一个在 Hadoop 平台上实现的 PB 级图挖掘库。Pegasus 的核心贡献是提出了 GIM-V(Generalized Iterated Matrix-Vector Multiplication) 原语,将许多图挖掘操作(包括连通分量计算)统一为广义的迭代矩阵-向量乘法。

Pegasus 中的连通分量算法本质上是一个标签传播(Label Propagation)过程,也被称为 Hash-Min:每个节点维护一个标签(初始为自身 ID),每轮迭代中将自己当前的最小标签传播给所有邻居,邻居则取收到的所有标签中的最小值作为新标签。当所有标签稳定不变时,拥有相同标签的节点属于同一连通分量。

Hash-Min 的每轮通信量是线性的 O(V + E),但迭代轮数为 O(d),其中 d 是图的直径。对于直径较小的社交网络图(通常 d < 20),这个算法表现不错。但对于直径较大的图(如路径图或某些稀疏图),迭代轮数可能非常多,实际性能急剧恶化。

4.2 Hash-to-Min 与 Hash-Greater-to-Min(2013)

Rastogi, V., Machanavajjhala, A., Chitnis, L. & Das Sarma, A. (2013). “Finding Connected Components in Map-Reduce in Logarithmic Rounds.” IEEE 29th International Conference on Data Engineering (ICDE), 50–61。论文全文可在 arxiv.org/abs/1203.5387 获取。

2013 年,Rastogi、Machanavajjhala、Chitnis 和 Das Sarma 在 IEEE ICDE 上发表了一项重要工作,系统地分析了 MapReduce 上连通分量算法的设计空间。他们提出了一个统一的 Map-Reduce 算法框架(Algorithm 1),通过不同的哈希函数 h 和合并函数 m 来产生不同的算法变体。

关键的算法变体包括:

  • Hash-Min(即 Pegasus):每个节点仅将最小标签传播给邻居,通信量低但收敛慢,需要 O(d) 轮。
  • Hash-to-All:每个节点将完整的簇信息发送给簇内所有成员,收敛快(O(log d) 轮),但通信量高达 O(n·|V| + |E|),对大连通分量不可行。
  • Hash-to-Min:一个巧妙的折中方案——每个节点将完整簇信息仅发送给簇内最小节点的 Reducer,其他节点只收到最小节点的 ID。这大幅降低了通信量,同时保持了较快的收敛速度。
  • Hash-Greater-to-Min:对 Shiloach-Vishkin PRAM 算法的高效 MapReduce 移植,可证明在 3 log n 轮内收敛,通信量为 2(|V| + |E|)。这是第一个同时具有对数轮次和线性通信量的 MapReduce 连通分量算法。

下表总结了各算法的复杂度对比(引自该论文 Table I):

算法迭代轮数每轮通信量
Pegasus / Hash-MinO(d)O(|V| + |E|)
Hash-to-AllO(log d)O(n·|V| + |E|)
Hash-Greater-to-Min3 log n2(|V| + |E|)

其中 n 是最大连通分量的节点数,d 是图的直径。

4.3 CC-MR(2012)

Seidl, T., Boden, B. & Fries, S. (2012). “CC-MR – Finding Connected Components in Huge Graphs with MapReduce.” ECML PKDD 2012, LNCS 7523, 458–473。

同一时期,Seidl、Boden 和 Fries 在 ECML PKDD 2012 上提出了 CC-MR 算法。CC-MR 在迭代轮数、通信开销和实际运行时间上都显著优于当时已有的 MapReduce 方法。它的创新之处在于利用了二次排序(Secondary Sorting)等 MapReduce 特性来优化中间数据的处理。

4.4 Alternating Algorithm:Large-Star + Small-Star(2014)

Kiveris, R., Lattanzi, S., Mirrokni, V., Rastogi, V. & Vassilvitskii, S. (2014). “Connected Components in MapReduce and Beyond.” Proceedings of the ACM Symposium on Cloud Computing (SoCC), 1–13。论文可在 ACM Digital Library 获取。

2014 年,来自 Google 的 Kiveris、Lattanzi、Mirrokni、Rastogi 和 Vassilvitskii 在 ACM SoCC(Symposium on Cloud Computing)上发表了论文 “Connected Components in MapReduce and Beyond”,提出了交替算法(Alternating Algorithm)。这篇论文是本文叙述的终点,也是 BigFunctions 中 connected_components 函数所采用的算法。

核心思想

与前面的算法类似,交替算法也是一种标签传播方法,但它引入了一个精妙的创新:在交替的方向上传播有向边,同时对节点标签执行最小值归约。具体来说,算法在每轮迭代中交替执行两个操作:

Large-Star(大星操作): 对于每个节点 v,找到它的所有邻居中标签比自己大的节点,将这些”大邻居”的标签重定向到 v 的标签(即让大的向小的靠拢)。直观地理解,这一步将每个节点周围标签较大的邻居都”拉”到自己身边,形成一个以较小标签为中心的星形结构。

Small-Star(小星操作): 对于每个节点 v,找到它的所有邻居中标签比自己小的节点,将 v 的标签重定向到最小邻居的标签(即让当前节点去追随更小的标签)。这一步让节点主动去寻找并跟随更小的标签。

为什么需要两个操作交替执行?

单独执行任何一个操作都不足以保证收敛。Large-Star 擅长将高标签节点拉向低标签节点,但可能遗漏某些传播方向;Small-Star 擅长让节点主动追随更小的标签,但也存在盲区。交替执行两者,标签可以在图中双向充分传播,确保同一连通分量内的所有节点最终收敛到同一个最小标签。

性能保证

  • 收敛速度: 已证明在 O(log²n) 轮内收敛,猜想(conjecture)在 O(log n) 轮内收敛。论文中的实验也支持对数级收敛的猜想。
  • 空间效率: 每轮迭代的数据量是线性的,不会出现某些算法(如 Hash-to-All)中的二次方数据膨胀问题。
  • 实际性能: 在论文的实验中,该算法比之前最好的算法快 3 到 15 倍(纯 MapReduce 实现),使用分布式哈希表(DHT)增强后更是快 10 到 30 倍,能够轻松扩展到拥有数千亿条边的图上。

为什么适合 BigQuery?

交替算法的每一步操作——“找邻居”、“取最小值”、“更新标签”——都可以自然地表达为 SQL 中的 JOIN + GROUP BY + MIN 操作。这使得它非常适合在 BigQuery 这样的 MPP(大规模并行处理)数据库中用纯 SQL 迭代实现,无需将数据导出到专门的图处理系统。这正是 BigFunctions 选择此算法的原因。

BigFunctions 文档中指出,该实现在每轮迭代中会持久化中间结果,整体成本大约相当于对输入表进行 15 到 30 次扫描。考虑到输入表只有两列(边的两个端点),这个成本通常是合理的。


五、后续发展(2014 年之后)

Google 2014 年论文之后,连通分量计算的研究并没有停止。后续的工作主要集中在以下几个方向:

5.1 进一步优化 MapReduce 实现

  • Cracker(Lulli et al., 2015):提出了一种新的分布式算法,通过将大图”碾碎”成小的连通片段来加速计算。
  • PACC(Park et al., 2020):提出了分区感知(Partition-Aware)的连通分量算法,通过两步处理(分区 + 计算)、边过滤和草图(sketching)三项技术,在实际数据集上比之前的最优算法快了最多 10.7 倍。
  • UniCon:将 Large-Star 和 Small-Star 合并为单一的 UniStar 操作,进一步简化算法并优化数据膨胀问题,成功处理了拥有 1290 亿条边的超大图。

5.2 面向 SQL/MPP 数据库的算法

2019 年,Bögeholz、Brand 和 Todor 发表了 “In-database Connected Component Analysis”,专门针对 MPP 关系数据库设计了 Randomised Contraction 算法。该算法的核心思想是:通过随机化的图收缩,每一步将图缩小到原来的一个常数比例,经过对数轮次后图缩小到只剩孤立点,这些孤立点就代表了各个连通分量。

与前述算法不同,Randomised Contraction 使用随机化来避免最坏情况,保证对任意输入图都能在期望 O(log |V|) 轮 SQL 查询内完成,且空间需求严格线性。实验表明,它在 MPP 数据库中的表现优于其他领先的连通分量算法。

参考文献: Bögeholz, H., Brand, M. & Todor, R.-A. (2019). “In-database Connected Component Analysis.” arXiv:1802.09478v2。论文可在 arxiv.org/abs/1802.09478 获取。

5.3 基于线性代数的分布式实现

Zhang、Azad 和 Buluç 基于 Shiloach-Vishkin 算法,设计了用稀疏线性代数原语实现的 LACCFastSV 算法,可以在 GraphBLAS 框架上高效运行。在 Cray XC40 超级计算机上,这些算法扩展到了 4096 个节点(262,144 个核心),处理超过 500 亿条边的图。

参考文献: Zhang, Y., Azad, A. & Buluç, A. (2020). “Parallel Algorithms for Finding Connected Components using Linear Algebra.” Journal of Parallel and Distributed Computing。论文可在 JPDC 获取。


六、演进脉络总结

下表总结了连通分量计算算法的关键里程碑:

年代计算模型代表算法/工作核心贡献
经典单机DFS/BFSO(V+E),最直观但无法并行
经典单机Union-Find (Tarjan, 1975)近常数均摊时间,支持流式处理
1982CRCW PRAMShiloach-VishkinO(log n) 并行算法,引入 Hooking + Shortcutting 范式
1991CRCW PRAMGazit最优随机化 PRAM 算法
1992EREW PRAMKarger-Nisan-Parnas在更严格的 EREW PRAM 模型上达到随机化 O(log n)
2009MapReducePegasus / Hash-Min首个 PB 级 MapReduce 图挖掘系统,O(d) 轮
2012MapReduceCC-MR优化通信开销和迭代次数
2013MapReduceHash-to-Min / Hash-Greater-to-Min首个对数轮+线性通信的 MapReduce 算法
2014MapReduceAlternating Algorithm (Kiveris et al.)Large-Star + Small-Star,O(log²n) 收敛,工业级性能
2015+多种Cracker, PACC, UniCon 等数据膨胀优化、负载均衡、超大规模图
2019SQL/MPPRandomised Contraction纯 SQL 实现,随机化保证 O(log n) 轮

七、启示与思考

回顾这段演进历程,有几个值得注意的趋势:

理论与实践的鸿沟。 PRAM 上的算法在 1980–1990 年代就已经达到了理论最优,但直到 2010 年代才被成功迁移到实际的分布式系统中。将共享内存算法转化为消息传递范式的分布式算法,远非简单的代码翻译。

简洁设计的力量。 Alternating Algorithm 的成功很大程度上来自于它的简洁性——两个操作交替执行,每个操作都可以用一条 SQL 查询表达。这种简洁性使得它能够在各种不同的计算平台(MapReduce、BSP、SQL 数据库)上轻松实现。

“够用”的收敛保证。 虽然 O(log²n) 的收敛界不如理论最优的 O(log n),但在实际应用中这个差距几乎不可感知。实验数据持续表明算法以 O(log n) 的速度收敛,这个猜想至今仍然开放。

数据基础设施决定了算法选择。 BigFunctions 选择 Alternating Algorithm 而非更新的 PACC 或 Randomised Contraction,很大程度上是因为 BigQuery 的 SQL 接口和计费模型与该算法的特性完美匹配。在算法设计中,计算模型和基础设施往往比理论上最优的渐近复杂度更重要。


参考文献汇总

  1. Shiloach, Y. & Vishkin, U. (1982). “An O(log n) Parallel Connectivity Algorithm.” Journal of Algorithms, 3(1), 57–67.
  2. Awerbuch, B. & Shiloach, Y. (1987). “New Connectivity and MSF Algorithms for Shuffle-Exchange Network and PRAM.” IEEE Transactions on Computers, 100(10), 1258–1263.
  3. Gazit, H. (1991). “An Optimal Randomized Parallel Algorithm for Finding Connected Components in a Graph.” SIAM J. Comput., 20(6), 1046–1067.
  4. Karger, D. R., Nisan, N. & Parnas, M. (1992). “Fast Connected Components Algorithms for the EREW PRAM.” SPAA ‘92, 373–381. 期刊版:SIAM J. Comput., 28(3), 1021–1034, 1999。
  5. Dean, J. & Ghemawat, S. (2004). “MapReduce: Simplified Data Processing on Large Clusters.” OSDI ‘04, 137–150.
  6. Kang, U., Tsourakakis, C. E. & Faloutsos, C. (2009). “PEGASUS: A Peta-Scale Graph Mining System — Implementation and Observations.” IEEE ICDM 2009.
  7. Seidl, T., Boden, B. & Fries, S. (2012). “CC-MR – Finding Connected Components in Huge Graphs with MapReduce.” ECML PKDD 2012, LNCS 7523, 458–473.
  8. Rastogi, V., Machanavajjhala, A., Chitnis, L. & Das Sarma, A. (2013). “Finding Connected Components in Map-Reduce in Logarithmic Rounds.” IEEE ICDE 2013, 50–61.
  9. Kiveris, R., Lattanzi, S., Mirrokni, V., Rastogi, V. & Vassilvitskii, S. (2014). “Connected Components in MapReduce and Beyond.” ACM SoCC 2014, 1–13.
  10. Lulli, A., Carlini, E., Dazzi, P., Lucchese, C. & Ricci, L. (2015). “Cracker: Crumbling Large Graphs into Connected Components.” IEEE ISCC 2015, 574–581.
  11. Bögeholz, H., Brand, M. & Todor, R.-A. (2019). “In-database Connected Component Analysis.” arXiv:1802.09478v2.
  12. Park, H.-M., Park, N., Myaeng, S.-H. & Kang, U. (2020). “PACC: Large Scale Connected Component Computation on Hadoop and Spark.” PLOS ONE, 15(3).
  13. Zhang, Y., Azad, A. & Buluç, A. (2020). “Parallel Algorithms for Finding Connected Components using Linear Algebra.” Journal of Parallel and Distributed Computing.
  14. BigFunctions connected_components 文档:https://unytics.io/bigfunctions/bigfunctions/connected_components/