Connected Component from DFS to BigQuery

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

Contents

去年迁移地图处理 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

这个图绘制出来就会如右图所示,一共有两个连通分量:{1, 2, 3, 4} 和 {5, 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 所在的集合合并。处理完所有边后,属于同一集合的节点就属于同一连通分量。

并查集支持两种核心操作: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

并查集的优势在于它可以按任意顺序处理边,甚至支持流式处理。但和 DFS/BFS 一样,它本质上仍是单机顺序算法,存在指针追逐问题:当前要访问的内存地址依赖于上一次读取出来的内存地址。所以也不足以处理大量数据。

基于并行计算 PRAM 模型的算法

想让成百上千个处理器同时干活,就必须先回答一个更基础的问题:多个处理器在同一时刻读写同一块内存时,究竟会发生什么? PRAM 并行随机存取机 (Parallel Random Access Machine) 就是为了回答这个问题而被造出来的模型。

PRAM 的基本假设

PRAM 模型由 Fortune, S. & Wyllie, J. (1978). “Parallelism in Random Access Machines.” STOC ‘78, 114–118 正式提出。系统的教材式介绍可参考 JáJá, J. (1992). An Introduction to Parallel Algorithms

PRAM 可以理解为把单机的 RAM (Random Access Machine) 模型原样复制若干份,再让它们共享同一块内存:

  • 处理器pp 个处理器 P1,,PpP_1, \dots, P_p,每个都是一台标准 RAM,拥有自己的局部寄存器,并且知道自己的编号 ii。很多算法就靠这个编号来分工,比如”第 ii 个处理器负责第 ii 条边”。
  • 共享内存:一块容量无限的全局内存,任何处理器都能在单位时间内访问其中任意一个单元。没有 cache,没有 NUMA,没有网络延迟,也没有带宽上限。
  • 同步执行:所有处理器步调一致(lock-step),每一步分为三个阶段——从共享内存读取、本地计算、向共享内存写回。所有处理器完成当前这一步之后,才会一起进入下一步。

可以看出这个模型极其慷慨:它把”通信”的成本直接抹成了零,只留下”计算需要多少轮”这一个变量。这当然不现实,但好处是让算法设计者可以先专心思考并行性本身的极限:在一个完全没有通信开销的理想世界里,一个问题最少需要几轮才能算完?

Brent 定理:一个需要 TT 轮、总工作量为 WW 的并行算法,可以在 pp 个处理器上用 O(W/p+T)O(W/p + T) 的时间模拟。这意味着处理器不够用时也不会吃太大的亏,因此算法设计者可以放心地先假设”处理器管够”。参见 Brent, R. P. (1974). “The Parallel Evaluation of General Arithmetic Expressions.” Journal of the ACM, 21(2), 201–206。

为此 PRAM 有两个评价指标:时间 TT(轮数)和处理器数 pp。两者的乘积 W=TpW = T \cdot p 称为工作量(work),代表所有处理器忙碌时间的总和。如果一个并行算法的 work 与最好的串行算法同阶,就称它是工作最优(work-optimal)的。连通分量的串行代价是 O(V+E)O(V + E),所以”用 O(V+E)O(V + E) 个处理器跑 O(logn)O(\log n) 轮”虽然够快,但 work 达到了 O((V+E)logn)O((V + E)\log n),比串行多做了 logn\log n 倍的无用功——这正是后来 Gazit 等人要改进的地方。

读写冲突:EREW、CREW 与 CRCW

模型里唯一需要认真定义的细节是:当多个处理器在同一步里访问同一个内存单元时,该怎么办? 对”并发读”和”并发写”分别选择允许或禁止,就得到了 PRAM 的几个变体:

模型全称同时读同时写
EREWExclusive Read Exclusive Write禁止禁止
CREWConcurrent Read Exclusive Write允许禁止
CRCWConcurrent Read Concurrent Write允许允许

并发读还算好理解,无非是把同一个值广播给多个处理器 理论上还存在 ERCW(禁止并发读但允许并发写)这一格,但既然写冲突这个更难的问题都解决了,再去限制读就没什么意义,因此几乎无人研究。 ;真正麻烦的是并发写:如果 100 个处理器同时往同一个单元写入 100 个不同的值,最后留下的应该是哪一个?于是 CRCW 内部又分成了几个流派:

  • Common(一致写):只有当所有写入者要写的值完全相同时才合法,否则就算算法本身有 bug。
  • Arbitrary(任意写):随便谁赢都行,算法必须保证无论谁赢结果都正确。
  • Priority(优先写):编号最小的处理器赢。
  • Combining(归约写):由硬件把所有写入值做一次归约(取 min / max / 求和)之后再写入。

这几种模型的能力依次递增:EREWCREWCRCW(Common)ArbitraryPriority\text{EREW} \le \text{CREW} \le \text{CRCW(Common)} \le \text{Arbitrary} \le \text{Priority}。不过差距并没有想象中那么大——任何 CRCW 算法都可以在 EREW PRAM 上模拟,代价是 O(logp)O(\log p) 的减速,做法是用排序和前缀和把并发访问”串行化”掉。

正因如此,在 CRCW 上做出 O(logn)O(\log n) 和在 EREW 上做出 O(logn)O(\log n) 是两件难度完全不同的事情:前者是 1982 年 Shiloach 和 Vishkin 的成果,后者要等到十年之后 Karger 等人的随机化算法。这也是为什么后文的表格里,每个算法都要专门标注它所依赖的模型。

而对连通分量来说,CRCW 尤其自然:算法的核心动作是”把邻居们的标签取最小值”,在 Combining CRCW 上这直接就是一条并发 min-write 指令;即便退到 Arbitrary CRCW,“任意一个邻居胜出”通常也只影响收敛的轮数,不影响最终结果的正确性。这就是这一时期的算法几乎都默认站在 CRCW 上的原因。

Shiloach-Vishkin 算法(1982)

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

1982 年,Yossi Shiloach 和 Uzi Vishkin 在 Journal of Algorithms 上发表了一个具有开创性意义的并行连通分量算法。它在 CRCW PRAM 上运行,O(logn)O(\log n) 时间,O(V+E)O(V + E) 个处理器。

数据结构:父指针森林

算法的全部状态就是一个数组 PPP[v]P[v] 是节点 vv 的父节点,根节点 rr 满足 P[r]=rP[r] = r。初始时 P[v]=vP[v] = v,每个节点自成一棵单点树。

这里有个容易被跳过但很关键的设计选择:为什么要维护一片树,而不是像并查集那样直接给每个节点存一个标签? 因为标签是”平”的,要让一个标签走到 kk 跳之外的节点,就得传播 kk 轮;而树是”有结构”的,一次指针跳跃就能把所有节点到根的距离减半O(logn)O(\log n) 这个结果,本质上就是从这个减半里来的。

算法维持的不变式是:PP 始终构成一片有根树的森林,同一棵树里的节点必定属于同一个连通分量。算法结束时每棵树都被压成深度为 1 的,根的编号就是这个分量的标签。

三个操作

Shiloach-Vishkin 算法的钩连与指针跳跃 一轮迭代:钩连让树变高,指针跳跃再把它压平。两者交替,森林在 O(log n) 轮内塌缩成星。

① 条件钩连(conditional hooking)。 对每条边 (u,v)(u, v) 并行地看:如果 uu 的父亲已经是一棵树的根,而 vv 那边的父亲编号更小,就把 uu 这棵树的根挂到 vv 的父亲下面。用一句话说就是——一条边如果连着两棵不同的树,就让编号大的那棵认编号小的那棵当爹

这一步正是需要 CRCW 的地方:可能有很多条边同时想改写同一个 P[root]P[\text{root}]。在 Arbitrary CRCW 下随便谁赢都行,因为无论谁赢,写进去的值都比原来的根小,正确性不受影响,只影响收敛的快慢。

② 指针跳跃(shortcutting)。 所有节点并行执行 P[v]P[P[v]]P[v] \leftarrow P[P[v]],也就是”直接认祖父当爹”。每执行一次,每个节点到根的距离减半,树高从 hh 掉到 h/2\lceil h/2 \rceil

③ 停滞钩连(stagnant hooking)。 只有①还不够。设想一棵树的根 rr,它周围所有的边都通向编号比它大的树——条件"P[v]<P[u]P[v] < P[u]"处处不成立,于是这一轮谁也没动,算法卡住了,但活儿还没干完。所以 SV 补了第二次钩连:对于在①里没能挂出去的根,允许它无条件地挂到任意一个邻居的树上。

// P[v] 是 v 的父节点;根节点满足 P[r] = r
parallel for v ∈ V:
    P[v] ← v

repeat until P 不再变化:

    // ① 条件钩连:把编号大的根挂到编号小的那边
    parallel for (u, v) ∈ E:
        if P[P[u]] = P[u] and P[v] < P[u]:
            P[P[u]] ← P[v]           // 并发写:Arbitrary CRCW,谁赢都对

    // ② 停滞钩连:给①里谁也没挂上的根一次无条件的机会
    parallel for (u, v) ∈ E:
        if P[P[u]] = P[u] and P[u] 在①中未被写过:
            P[P[u]] ← P[v]

    // ③ 指针跳跃:直接认祖父,树高减半
    parallel for v ∈ V:
        P[v] ← P[P[v]]

右图就是一轮迭代的样子。初始有两棵树,根分别是 2 和 3,中间隔着一条图上的边 (5,4)(5, 4)。钩连这一步,边 (4,5)(4, 5) 让 3 这棵树的根挂到了 2 上——注意此时树变深了,4 和 9 离根有 2 跳。紧接着的指针跳跃把这个深度重新压掉:4 和 9 原本指向 3,而 P[3]P[3] 已经是 2,于是它们一步就直接指向了 2。

一高一压,正是这个算法的全部节奏。 钩连负责合并,它的副作用是把树拉长;指针跳跃负责把拉长的部分收回来,让下一轮的钩连仍然发生在”根与根”之间。原文用一个关于森林结构的势函数证明了:无论输入图长什么样,O(logn)O(\log n) 轮之后一定收敛。

和并查集的关系

看到这里应该会觉得眼熟:SV 其实就是并查集的并行版本

并查集(串行)Shiloach-Vishkin(并行)
Union:合并两个集合Hooking:把一棵树挂到另一棵树上
路径压缩:Find 时顺手把路径拍平Shortcutting:所有节点同时认祖父
一条边一条边地处理所有边同时处理

差别全在最后一行。串行地处理边,天然不会有冲突;一次性并行处理所有边,就必须额外发明规则来避免成环和停滞——①和③这两条别扭的规则,全都是为并行付的税。

后续的 PRAM 改进

Awerbuch, B. & Shiloach, Y. (1987). “New Connectivity and MSF Algorithms for Shuffle-Exchange Network and PRAM.” IEEE Transactions on Computers, C-36(10), 1258–1263。

Johnson, D. B. & Metaxas, P. (1991). “Connected Components in O(log3/2V)O(\log^{3/2}|V|) Parallel Time for the CREW PRAM.” FOCS ‘91, 688–697。

SV 之后的十年,PRAM 上的改进基本沿着两条线走:把 work 降到最优,以及在更弱的模型上做到同样快

  • Awerbuch & Shiloach (1987) 简化了钩连的判定规则,不再需要显式记录”谁在①里被写过”,把三步压成了一套统一的规则,实现起来简单得多。后来几乎所有工程实现参照的都是这个版本。

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。期刊版:SIAM J. Comput., 28(3), 1021–1034, 1999。

  • Gazit (1991) 解决了 work 的问题。SV 用 O(V+E)O(V+E) 个处理器跑 O(logn)O(\log n) 轮,work 是 O((V+E)logn)O((V+E)\log n),比串行多做了 logn\log n 倍。Gazit 用随机采样先把图收缩一轮,做到 O(logn)O(\log n) 时间 + O((V+E)/logn)O((V + E)/\log n) 个处理器,work 与串行同阶——这是工作最优的。
  • Johnson & Metaxas (1991) 在禁止并发写的 CREW 上给出了 O(log3/2n)O(\log^{3/2} n) 的确定性算法。
  • Karger, Nisan & Parnas (1992/1999) 啃下了最硬的 EREW:连”并发读”都不许,min-write 这种一步到位的动作必须靠显式构造广播树来模拟。他们给出了随机化 O(logn)O(\log n) 和确定性 O(log1.5n)O(\log^{1.5} n) 的算法。

Eppstein, D. & Galil, Z. (1988). “Parallel Algorithmic Techniques for Combinatorial Computation.” Annual Review of Computer Science, 3, 233–283。

然而这些成果几乎都停留在纸面上。PRAM 模型”常被理论计算机科学家使用,却较少被实际并行机器的构建者采用”。

从 PRAM 到 MapReduce:模型是怎么迁移的

回头看 PRAM 的三条基本假设,放到一个真实的集群上,没有一条站得住:

“任意地址、单位时间” 一次本地内存访问大约 100 纳秒,而一次跨机器的网络往返在几百微秒到毫秒量级,差了三到四个数量级,这是无法忽视的性能差距。

“lock-step 同步” PRAM 假设所有处理器齐步走,而现实里做一次全局同步本身就是一笔可观的开销,何况集群里总有掉队的慢节点。

“没有失败” PRAM 里压根没有”机器挂掉”这个概念。而在几千台机器的大集群上,任务失败是每分钟都在发生的常态。

Valiant, L. G. (1990). “A Bridging Model for Parallel Computation.” Communications of the ACM, 33(8), 103–111。Valiant 用的词是 “bridging model”,一座架在硬件和算法之间的桥。

第一次认真讨论通信代价的,是 Valiant 在 1990 年提出的 BSP(Bulk Synchronous Parallel)。BSP 把计算切成一连串 “SuperStep”,每个 SuperStep 依次做三件事:本地计算、收发消息、全局同步。

一次 SuperStep 的代价为:

w+gh+Lw + g \cdot h + L

其中 ww 是本地计算量,hh 是这一步收发的消息条数,gg 是网络的单位消息成本,LL 是同步的固定开销。Google 后来的图计算系统 Pregel,用的就是 BSP。

Dean, J. & Ghemawat, S. (2004). “MapReduce: Simplified Data Processing on Large Clusters.” OSDI ‘04, 137–150。

而在 2004 年横空出世的 MapReduce 比 BSP 还要激进,它干脆把通信原语砍到只剩一个:按 key 分组

map:      (k, v)      →  list(k', v')
shuffle:  按 k' 把所有记录分到同一个 reducer          ← 唯一的通信原语
reduce:   (k', [v'])  →  输出

这个设计看似限制了灵活性,但却给使用者带来了极大的便利:通信模式是框架规定死的,框架就能自动完成分区、调度、失败重试和落盘,写算法的人不用再管这些细枝末节的东西了。

从 PRAM 到 BSP 再到 MapReduce 的计算模型演进
三代并行模型的演进。指针跳跃 P[P[v]] 在 PRAM 里是一次免费的内存访问,到了 MapReduce 里就是一次 JOIN。

Karloff, H., Suri, S. & Vassilvitskii, S. (2010). “A Model of Computation for MapReduce.” SODA ‘10, 938–948。文中的 Vassilvitskii 也是四年后那篇交替算法论文的作者之一。

那么 PRAM 上积累了三十年的算法,能不能搬到 MapReduce 上呢?可以。2010 年提出的 MRC 模型 解决了这个问题。

为了把 MapReduce 形式化,加了三条约束:单台机器的内存是 O(n1ϵ)O(n^{1-\epsilon})(装不下整个输入)、机器台数是 O(n1ϵ)O(n^{1-\epsilon})、以及同一个 key 对应的所有数据必须能装进一台机器。跑 O(login)O(\log^i n) 轮的问题被归入 MRCi\mathcal{MRC}^i 类。

任何使用 O(n22ϵ)O(n^{2-2\epsilon}) 总内存、跑 tt 步的 CREW PRAM 算法,都可以用 O(t)O(t) 轮 MapReduce 模拟出来。

这意味着 PRAM 上的 O(logn)O(\log n) 结果,理论上都可以自动转换成 O(logn)O(\log n) 轮的 MapReduce 算法。

但”能模拟”和”划算”是两码事。一轮 MapReduce 的固定开销包括作业启动、一次全量 shuffle、以及把中间结果写回磁盘(在 Hadoop 时代还得写三副本)。所以对于 MapReduce 来说,每次 Shuffle 的代价都不低,尤其是需要交换大量数据的时候。而每轮要搬动多少数据,这是在 PRAM 模型上完全不会考虑的问题。这就导致很多理论上在 PRAM 上效果很好的算法,到了 MapReduce 上就不灵。

PRAM 里的动作MapReduce 里的对应代价
P[v]P[v]map 端读一条记录几乎免费
并发 min-writereduce 端 GROUP BY + MIN搭一轮作业的便车
指针跳跃 P[P[v]]P[P[v]]自连接 JOIN一整轮作业
lock-step 同步作业结束时的隐式 barrier启动 + 落盘开销

JOIN 的代价非常昂贵,但又无法避免。所以之后的算法本质上都在琢磨同一个问题:怎样用更少的 JOIN,做完同样多的指针跳跃。

MapReduce 时代的连通分量算法

这个时候算法的形态基本都一模一样:图就是一张边表,算法就是一个”把边表变成边表”的作业,反复跑到不动点为止。区别只在于每一轮中怎么改写这张表,以及要跑多少轮。

Pegasus / Hash-Min

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 年的 Pegasus 是第一个跑在 Hadoop 上的 PB 级图挖掘库。它的核心洞见是:PageRank、直径估计、连通分量这些图算法其实都是在反复做矩阵-向量乘法。

于是他们把这件事抽象成一个原语 GIM-V(Generalized Iterated Matrix-Vector Multiplication):

把标准矩阵乘法 jMijvj\sum_j M_{ij} v_j 里的”乘”和”加”换成两个可自定义的函数,再加一个把结果写回的 assign 函数。连通分量就是在 (min,×)(\min, \times) 这个半环上做迭代:

vinew=min(vi, minjN(i)vj)v_i^{\text{new}} = \min\Big(v_i,\ \min_{j \in N(i)} v_j\Big)

翻译成 MapReduce 就是两行:

// 状态:每个节点 v 的当前标签 L[v],初始 L[v] = v
map(v, L[v], N(v)):
    emit(v, L[v])                      // 保留自己的标签
    for u ∈ N(v): emit(u, L[v])        // 把标签发给每个邻居

reduce(v, labels):
    L[v] ← min(labels)                 // 取收到的最小值

这就是后来被称为 Hash-Min 的算法:每轮把自己的标签发给所有邻居,再从收到的标签里取最小值。每轮的通信量是干净的 O(V+E)O(V + E)

但问题出在轮数上。标签每轮只能前进一跳

Hash-Min 在路径图上的执行过程,标签每轮前进一跳
最小标签 1 像涨潮一样,一轮往前推一格。图的直径有多大,就要跑多少轮。

所以 Hash-Min 算法的轮数是 O(d)O(d)dd 是图的直径。

如果是计算社交网络来说这没什么问题。按照六度分隔理论,我隔着 6 个人就能找到世界上任何人,所以 dd 撑死也就 20,基本上跑二十来轮就收敛了。

但换成道路网就不太行,因为道路就是要方便大家前往任何需要的地方。道路的连通分量图是极度稀疏的近平面图,一条国道可能串起几万个路口,直径轻轻松松上万。在这种图上跑 Hash-Min,等于要提交上万轮 MapReduce 作业,而且每轮都是一次全量 shuffle 加落盘。这个算法在这种场景里是基本没发用的。

Rastogi 等人的统一框架

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 等人设计了一套新的框架:

每个节点 vv 维护一个簇 CvC_v(一个节点集合,初始是 vv 和它的邻居)。每一轮迭代做两件事,先用一个哈希函数 hh 决定“把 CvC_v 的哪些部分发给谁”,再用一个合并函数 mm 把每个节点收到的东西并起来。换一对 (h,m)(h, m),就得到一个不同的算法。

Hash-to-All 与 Hash-to-Min 的消息模式对比
同样一个 5 个节点的簇:左边 20 条消息,右边 8 条。Hash-to-Min 让最小节点当集散中心,把 n² 降到 n。
  • Hash-Min:只把 min(Cv)\min(C_v) 发给邻居。通信量最省,O(d)O(d) 轮。
  • Hash-to-All:把整个 CvC_v 发给 CvC_v 里的每一个成员。这样每个节点每轮都能看到”邻居的邻居的全部”,簇的半径每轮翻倍,于是只要 O(logd)O(\log d) 轮。代价是每轮通信量高达 O(nV+E)O(n \cdot |V| + |E|),簇一大就直接爆炸。
  • Hash-to-Min:二者的折中。把整个 CvC_v 只发给簇内编号最小的那个节点,簇的合并在最小节点处一次性完成,其余成员只收到一个 ID。论文证明了它在路径图上 O(logn)O(\log n) 轮收敛,一般图上则是猜想。
  • Hash-Greater-to-Min:把簇只发给编号大于自己的那部分成员。这实际上是 SV 算法的一次高效 MapReduce 移植,并且可以严格证明 3logn3\log n 轮收敛、每轮通信量 2(V+E)2(|V| + |E|)这是第一个同时做到对数轮数和线性通信量的 MapReduce 连通分量算法。
算法迭代轮数每轮通信量
Hash-Min(即 Pegasus)O(d)O(|V| + |E|)
Hash-to-AllO(log d)O(n·|V| + |E|)
Hash-to-Min路径图 O(log n),一般图为猜想实测接近线性
Hash-Greater-to-Min3 log n2(|V| + |E|)

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

CC-MR

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 提出的 CC-MR 走的是另一条路:节点中心。每个 reducer 拿到一个节点 vv 和它当前的邻接表 N(v)N(v),令 m=min(N(v){v})m = \min(N(v) \cup \{v\}),然后把邻域里所有比 mm 大的节点统统改挂到 mm 上。也就是说,每一轮每个节点都在自己的邻域里做一次”局部星形化”。

它真正的巧思在工程上:CC-MR 利用了 MapReduce 的二次排序(secondary sorting),让 reducer 拿到的 value 迭代器第一个元素就是最小的邻居。这样就不必把整个邻接表读进内存再求最小值,可以完全流式地处理。在幂律分布的真实图里,总有那么几个度数几百万的超级节点,能不能流式处理它们,直接决定了作业会不会 OOM。

现在回头看,CC-MR 其实已经非常接近下面的算法了。“把邻域里比最小值大的节点全改挂到最小值”,这句话几乎就是 Large-Star 的定义。差的那一步是:Kiveris 等人把这个动作拆成了两个方向,并且给出了收敛性证明。

伟大时刻:Alternating Algorithm(Large-Star + Small-Star)

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 发表了 “Connected Components in MapReduce and Beyond”,提出了交替算法(Alternating Algorithm)。

这个算法的全部状态是一个有向边集,(u,v)(u, v) 读作”uu 指向 vv“。初始时把无向图的每条边双向展开放进去。算法每一轮把这个边集重写一遍,最终收敛到”每个节点直接指向自己所在分量的最小编号”。

对每个节点 vv,记它当前的邻居集合为 N(v)N(v),记 mv=min(N(v){v})m_v = \min\big(N(v) \cup \{v\}\big),也就是邻域(含自己)里的最小编号:

  • Large-Star(vv):N(v)N(v) 中所有严格大于 vv 的邻居改挂到 mvm_v
  • Small-Star(vv):N(v){v}N(v) \cup \{v\} 中所有不大于 vv 的邻居(包括 vv 自己)改挂到 mvm_v

用伪代码来描述就是:

map(边 (u, v)):
    emit(u, v);  emit(v, u)                    // 双向展开成邻接表

reduce(v, N(v)):                               // Large-Star
    m ← min(N(v) ∪ {v})
    for u ∈ N(v) where u > v:      输出边 (u, m)

reduce(v, N(v)):                               // Small-Star
    m ← min(N(v) ∪ {v})
    for u ∈ (N(v) ∪ {v}) where u ≤ v:  输出边 (u, m)

每个操作恰好是一轮 MapReduce:map 端双向展开,reduce 端按节点分组、求 min、重写边。没有任何多余的 JOIN。

Large-Star 与 Small-Star 在一个小图上的完整执行过程 两轮交替:Large-Star 让 5 同时够到 2 和 3,Small-Star 随即在 5 的邻域里把 3 并给 2。(自环边已省略。)

这里我有这样的一个图:节点 {2,3,4,5,7,9}\{2,3,4,5,7,9\},边 25, 27, 34, 39, 452\text{–}5,\ 2\text{–}7,\ 3\text{–}4,\ 3\text{–}9,\ 4\text{–}5。图里有两个”局部最小值” 2 和 3,它们之间隔着 454\text{–}5 这条边。

第一轮 Large-Star。v=4v = 4 这个 reducer 里,N(4){4}={3,4,5}N(4) \cup \{4\} = \{3,4,5\},最小是 3,而 5 比 4 大,于是边 (5,4)(5, 4) 被改写成了 (5,3)(5, 3)。与此同时 v=2v=2 那边把 5 和 7 挂到了 2 上。现在节点 5 现在同时指向 2 和 3。两个原本隔了两跳的局部最小值,被拉到了同一个节点的邻域里。

第一轮 Small-Star。 现在轮到 v=5v = 5 这个 reducer:N(5){5}={2,3,5}N(5) \cup \{5\} = \{2,3,5\},最小是 2,而 3 不大于 5,于是 3 被挂到了 2 下面。两个局部最小值在这里合并了。

第二轮 Large-Star 把剩下的 4 和 9 从 3 那里直接拉到 2,边集不再变化,收敛。总共两轮。

也就是两个刚好互补的操作交替执行:Large-Star 负责把大节点拉到局部最小值身边,顺带制造出”同时够到两个局部最小值”的节点;Small-Star 紧接着在这些节点的邻域里,把两个局部最小值合并掉。

如果只做 Large-Star,第一轮之后边集变成 {(4,3),(5,2),(5,3),(7,2),(9,3)}\{(4,3), (5,2), (5,3), (7,2), (9,3)\},然后就卡死了, 2 和 3 仍然是两个根。因为 Large-Star 只改写”比 vv 大的邻居”,而一个局部最小值自己的指针永远不会在自己的邻域里被改动,于是两个局部最小值之间谁也不肯让步。

如果只做 Small-Star,会停在 {(2,2),(3,2),(3,3),(4,2),(5,2),(7,2),(9,3)}\{(2,2), (3,2), (3,3), (4,2), (5,2), (7,2), (9,3)\}:9 还挂在 3 上,3 又挂在 2 上,树没被压平,也不是最终答案。因为 Small-Star 只改写”不大于 vv 的邻居”,它擅长把局部最小值并起来,却没本事把远处的大节点拉过来。

这个”一拉一并”的分工,和 Shiloach-Vishkin 算法里“钩连拉长、跳跃压平”其实差不多。三十二年过去,相同的思路在 MapReduce 模型上焕发新生。

论文证明了交替算法在 O(log2n)O(\log^2 n) 轮内收敛,并猜想真实的界是 O(logn)O(\log n)。这就相当快了。简单用路径图验证:

路径长度 nnHash-Min 轮数交替算法轮数
873
16154
32315
64636
1281277

交替算法轮数几乎就是 log2n\log_2 n,但对此 Hash-Min 就拉了。

还有一个交换数据量的问题。Large-Star 之后边集里会出现重复边,需要去重。此外他们还用 DHT 把已经收敛的小分量直接摘出去,只对大分量继续迭代。最终 MapReduce 实现比当时最好的算法快 3 到 15 倍,DHT 增强版快 10 到 30 倍,能扛住数千亿条边的图。

后续发展(2014 年之后)

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

进一步优化 MapReduce 实现

Lulli, A., Carlini, E., Dazzi, P., Lucchese, C. & Ricci, L. (2015). “Cracker: Crumbling Large Graphs into Connected Components.” IEEE ISCC 2015, 574–581。

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), e0229936。

Kim, C., Han, C. & Park, H.-M. (2022). “UniCon: A Unified Star-Operation to Efficiently Find Connected Components on a Cluster of Commodity Hardware.” PLOS ONE, 17(11), e0277527。

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

面向 SQL/MPP 数据库的算法

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

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

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

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

Zhang, Y., Azad, A. & Buluç, A. (2020). “Parallel Algorithms for Finding Connected Components Using Linear Algebra.” Journal of Parallel and Distributed Computing, 144, 14–27。论文可在 ScienceDirect 获取。

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

BigQuery 上的实现

BigFunctions 的 connected_components 文档:unytics.io/bigfunctions/bigfunctions/connected_components/。 GraphFrames 把这个实现叫作 “Two Phase” 或 “big star / small star”,在文档里直接引用了 Kiveris 等人这篇论文。它同时还保留了一个 Pregel 版本的实现,那个就是每轮推进一跳、轮数等于图直径的 Hash-Min。

按照 GraphFrames 文档的说法,我用的那个 g.connectedComponents.run() 的默认实现正是 Large-Star / Small-Star 交替算法。为什么不用更新的 PACC 或 Randomised Contraction 呢?这很大程度上是因为 BigQuery 的 SQL 接口和计费模型与该算法的特性完美匹配。

比如说 “找邻居”、“取最小值”、“改挂边” 的操作分别对应 SQL 里的 JOINGROUP BY ... MINSELECT。一轮 Large-Star 写出来是这样(示意,省略了终止判断):

-- edges(src, dst) 是当前的有向边集
WITH nbr AS (                                   -- 双向展开成邻接表
  SELECT src AS v, dst AS u FROM edges
  UNION ALL
  SELECT dst AS v, src AS u FROM edges
),
mins AS (                                       -- 每个节点邻域(含自己)的最小值 m
  SELECT v, LEAST(MIN(u), v) AS m
  FROM nbr GROUP BY v
)
SELECT DISTINCT n.u AS src, mins.m AS dst       -- 把更大的邻居改挂到 m
FROM nbr AS n JOIN mins ON n.v = mins.v
WHERE n.u > n.v

Small-Star 只需要把最后的 n.u > n.v 换成 n.u <= n.v,并把 vv 自己也加进邻域。整个算法就是这两段 SQL 交替跑到边集不再变化为止。BigFunctions 的文档提到,它的实现每轮会持久化中间结果,整体成本大约相当于对输入表做 15 到 30 次扫描。考虑到输入表只有两列,这个代价通常是可以接受的。


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

年代计算模型代表算法/工作核心贡献
经典单机DFS/BFSO(V+E),最直观但无法并行
经典单机Union-Find (Tarjan, 1975)近常数均摊时间,支持流式处理
1982CRCW PRAMShiloach-VishkinO(log n) 并行算法,引入 Hooking + Shortcutting 范式
1990BSP (Valiant)第一次给通信明码标价:w + g·h + L
1991CRCW PRAMGazit最优随机化 PRAM 算法
1992EREW PRAMKarger-Nisan-Parnas在更严格的 EREW PRAM 模型上达到随机化 O(log n)
2004MapReduce (Dean & Ghemawat)通信原语砍到只剩「按 key 分组」,换来自动容错
2009MapReducePegasus / Hash-Min首个 PB 级 MapReduce 图挖掘系统,O(d) 轮
2010MRC (Karloff-Suri-Vassilvitskii)证明 CREW PRAM 可被 O(t) 轮 MapReduce 模拟
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) 轮