Pregel
Pregel 是一个面向大规模图的分布式计算框架,把程序组织成顶点中心(vertex-centric)、按超步同步的形状。它要处理的图在规模上很极端:顶点可达数十亿、边可达数万亿,典型的例子是 Web 图与各种社交网络。
它同时也是"专用框架"这条路线的一个标志性样本 —— RDD 那篇把它当作论据之一,说明图计算可以只用 200 行库在通用抽象上表达出来。所以这一篇值得读两层:Pregel 自己定下的模型,以及它为此付出的专用化代价。
为什么图计算需要专门的框架
图算法有三个共同的形状特征,合起来让通用框架很难伺候:
- 内存访问的局部性很差;
- 每个顶点上的工作量极少;
- 执行过程中并行度会变。
把这套东西分发到多台机器上会放大局部性问题,同时提高"计算期间某台机器挂掉"的概率。
而实现一个处理大图的算法,当时只有四条路,每条都不合用:
| 可选路径 | 问题 |
|---|---|
| 自己造一套分布式基础设施 | 工作量巨大,而且每换一个算法或换一种图表示就得重来一遍 |
| 借现成的分布式平台 | 常常不适配图处理。MapReduce 用在很多大规模计算上很合适,也被用来挖大图,但会带来欠优的性能与可用性问题;它被扩展出过聚合、类 SQL 查询等能力,但这些扩展对更适合消息传递模型的图算法并不理想 |
| 用单机图算法库(BGL、LEDA、NetworkX、JDSL、Stanford GraphBase、FGL) | 限制了问题规模 |
| 用已有的并行图系统(Parallel BGL、CGMgraph) | 处理了并行图算法,但不解决容错,以及大规模分布式环境里其他要紧的问题 |
Pregel 的目标是补上这一格:一个可扩展、可容错、且 API 足以表达任意图算法的平台。
高层组织方式借鉴了 Valiant 的 BSP(Bulk Synchronous Parallel) 模型。
模型:超步与顶点状态机
输入是一个有向图:每个顶点由一个字符串顶点标识符唯一确定,并带一个用户定义、可修改的值;有向边从属于它的源顶点,每条边由一个可修改的值加一个目标顶点标识符构成。
计算过程是:输入(初始化图)→ 一串超步(superstep),超步之间由全局同步点分隔 → 输出,直到算法终止。
在一个超步内,所有顶点并行计算,每个都执行同一个用户定义函数。这个函数只负责描述**"顶点
- 读上一个超步发给
的消息; - 发消息给其他顶点(在下一个超步被收到);
- 修改
自己的状态,以及它的出边的状态; - 改图的拓扑(增删点边)。
边不是一等公民 —— 边没有关联的计算,只有值。
终止条件:投票停机
算法终止靠每个顶点投票停机(vote to halt):
规则是三句话:
- 超步 0 里每个顶点都处于 active,所有 active 顶点参与任一超步的计算;
- 顶点通过投票停机来自我停用 —— 含义是"除非被外部触发,我没有别的工作要做",此后框架不再执行它,除非它收到消息;
- 被消息重新唤醒的顶点必须显式地再次停用自己。
算法整体终止的条件是:所有顶点同时处于 inactive,且没有消息在途。
输出是顶点显式输出的值的集合。它常常是同构于输入的有向图,但这不是系统必需的性质 —— 计算过程中顶点和边都可以增删。聚类算法可能产出一小组从大图里选出的不连通顶点;图挖掘算法可能只输出聚合统计量。
一个最小例子:给定一张强连通图、每个顶点持有一个值,把最大值传播到每个顶点。每个超步里,任何从消息中学到更大值的顶点把它发给所有邻居;当某个超步里不再有顶点改变时算法终止。
为什么是纯消息传递
Pregel 只做消息传递,不做远程读、也不提供其他模拟共享内存的手段。两个理由:
- 表达力足够 —— 消息传递已经够强,不需要远程读;没有找到任何消息传递表达不了的图算法。
- 性能更好 —— 在集群环境里从远端机器读一个值会引入无法轻易隐藏的高延迟;而消息传递模型可以把消息异步地批量投递,从而摊薄延迟。
与"把图算法写成一串链式 MapReduce 调用"相比,差别在于:
- Pregel 让顶点和边留在做计算的那台机器上,网络只用于传消息;
- MapReduce 本质上是函数式的,所以把图算法写成链式 MapReduce 要求把图的整个状态从一个阶段传到下一个阶段 —— 一般会带来多得多的通信与随之而来的序列化开销;此外协调链式 MapReduce 的各步骤额外引入编程复杂度,而 Pregel 用超步迭代把这块消掉了。
表达的边界:顶点中心能写什么、写不了什么
"消息传递足以表达任意图算法"这句话,口径要收窄三层才能用:能表达(存在一份等价的顶点中心程序)、写得顺(不必为了绕开模型的缺失去搭额外结构)、跑得值(不为绕开缺陷付出成倍的通信或超步)。三层是三个不同的问题,把第一层当成全部,会在选型时误判。
顶点中心的形状约束只有一条,但很硬:Compute() 能看到的全部输入,就是"上一超步到达的消息 + 本顶点的值 + 本顶点的出边"。凡是这三样之外的信息,都得先想办法搬进来;搬运方式决定写起来顺不顺:
| 需要的东西 | 顶点中心里怎么拿到 | 代价 |
|---|---|---|
| 全局统计量(总数、最值、直方图) | Aggregator —— 本超步供给、下一超步全局可见 | 只提供归约,拿不到全部值,且慢一个超步 |
| 远方某个顶点的当前值 | 发一条请求消息,等对方在下一超步回复 | 往返两个超步,而回复到达时对方的值可能已经变了 |
| 自己的入边 | 拿不到 | 入边是源顶点所持列表的一部分,源顶点通常在另一台机器上 |
| 一份稳定的全图快照 | 拿不到 | 任何时刻看到的都是"上一超步结束时的值" |
后两行是同一个物理事实的两面:框架把顶点和它的出边放在同一台机器上,于是"出边可读、入边不可读";而"只看得到上一超步"是超步屏障的直接后果。这两条合起来划定了顶点中心的能力圈,也解释了为什么三处"看起来缺了什么"的地方都能被现有机制绕过去 —— 代价是慢一轮或换一种写法。
Aggregator 与"全局状态"的差别值得单独说清:它给的是一个归约结果,不是一份可读写的全局变量。想统计图的边总数,用 sum 聚合出度即可;但"想知道某个特定顶点当前的值"这类查询,聚合器帮不上忙 —— 只能走请求消息。这也解释了为什么可以放心省掉远程读:需要远端值时发消息、等一轮回复,通信本身可以批量投递;而远程读是一次不可批量的同步等待,延迟无从隐藏。
一个容易踩的写法陷阱是**"比较两个顶点的当前值"。顶点中心里没有"读一眼邻居 SendMessageTo(u, REQ) 再在下一超步取回复,
另一条边界落在确定性上:受限恢复要求用户算法确定;不确定的算法照样能跑,只是恢复时退回基本机制、重算量更大。随机化算法可以用"按超步与分区确定性播种伪随机数发生器"变成确定性的 —— 把算法改造成适配框架,这是常见做法而不是例外。
一段选型前的判据 —— 把算法改写成顶点中心之前,先问三个问题:
- 一个顶点做决策所需的输入,能不能全部从"上一超步的消息 + 自己的值 + 自己的出边"推出来?
- 算法里有没有"必须一次看到全图才能继续"的步骤?
- 有没有依赖"读到别人最新值"的地方?
三个都能答"能"或"没有",改写通常直接;只要有一个答不上,就得先设计"怎么把那条信息搬进消息",而这一步往往决定这个算法写起来是顺畅还是别扭。从 Pregel 自己的使用经验看,转到"像顶点一样思考"这个视角之后,API 是直觉、好用的 —— 这是模型层面的收获;但"能写"与"写得好"之间的距离,正是上面三问要提前量出来的。
把算法改造成顶点中心,常见的动作只有三个,落到写法上:
| 原算法里的东西 | 顶点中心里的写法 | 多付的代价 |
|---|---|---|
| 一个全局累积量 | Aggregator,每轮由各顶点供给、下一轮广播 | 慢一个超步;只能归约,不能任意读写 |
| 一次对远方数据的直接读取 | 请求消息 + 下一超步取回复 | 往返两个超步;回复到达时值可能已变 |
| 一处随机或不依赖输入的状态 | 把不确定性前移固化(播种、把外部状态当参数传进来) | 改写成本;换来受限恢复可用 |
第三类动作最容易被低估:把不确定性"前移固化"往往要改算法的形状,加一个参数解决不了。典型做法是"把随机决策提到超步 0、由分区号与超步号共同决定",这样任何一次重放都能得到同一串随机数。
一条附带的认识:顶点中心里"能表达"与"跑得起"通常差着超步数。凡是需要"看别人最新值"的地方,改写后都会变成"基于上一超步快照做判断",于是收敛轮数会变 —— 有时变多(要传的信息走得更慢),有时变少(不用等所有人)。这个差异在选型阶段就该算出来,等跑起来才发现轮数比预期多一倍,往往意味着算法选错了模型。
C++ API:每顶点一个值 + Compute()
写一个 Pregel 程序就是继承预定义的 Vertex 类。它的模板参数定义三个值类型:顶点值、边值、消息值。这种"每顶点/每边只有一个值"的统一性看起来受限,但用户可以用协议缓冲区这类能承载任意结构的类型。
template <typename VertexValue, typename EdgeValue, typename MessageValue>
class Vertex {
public:
virtual void Compute(MessageIterator* msgs) = 0;
const string& vertex_id() const;
int64 superstep() const;
const VertexValue& GetValue();
VertexValue* MutableValue();
OutEdgeIterator GetOutEdgeIterator();
void SendMessageTo(const string& dest_vertex, const MessageValue& message);
void VoteToHalt();
};用户覆盖 Compute(),它在每个 active 顶点、每个超步被执行。通过 GetValue() / MutableValue() 读写顶点值,通过出边迭代器读写出边的值。这些状态更新立即可见,但因为可见范围限于被修改的那个顶点,所以从不同顶点并发访问值不存在数据竞争。
一条设计上很关键的约束:顶点及其边的值是唯一跨超步存活的每顶点状态。把框架管理的图状态限制成"每顶点/每边一个值",简化了三件事 —— 主计算循环、图的分发、以及失败恢复。
消息传递
消息由一个消息值加一个目标顶点名构成,值是用户指定的模板参数。一个顶点在一个超步里可以发任意多条消息。超步 Compute() 时通过迭代器可见。
两条保证、一条不保证:
- 保证送达;
- 保证不重复;
- 迭代器里不保证消息顺序。
目标顶点不必是
任何消息的目标顶点不存在时,执行用户定义的 handler —— 例如创建缺失的顶点,或者把那条悬空边从源顶点上删掉。
Combiner:减少消息开销
发消息(尤其是发往另一台机器)是有开销的,用户可以在某些情况下帮忙降低它。举例:如果 Compute() 收到的是整数消息,而只有"和"有意义、各个具体值无所谓,那么系统就能把发给顶点
Combiner 默认不开启,原因是不存在机械的办法找出一个既有用、又与用户 Compute() 语义一致的合并函数 —— 这件事必须由用户声明。做法是继承 Combiner 类并覆盖 Combine()。
一条必须注意的边界:对"合并了哪些消息、分组方式、合并顺序"都没有任何保证,所以只能对可交换且可结合的操作启用 Combiner。
实测收益:在单源最短路(见下)这类算法上,消息流量减少到原来的四分之一以下。
Aggregator:全局通信与协调
Aggregator 是全局通信、监控与数据的机制:每个顶点在超步 min、max、sum,覆盖整数与字符串类型。
三类典型用法:
- 统计 —— 对每个顶点的出度做
sum就是图的边总数;更复杂的归约算子可以生成某个统计量的直方图; - 全局协调 —— 让
Compute()的一个分支先执行若干超步,直到一个and聚合器判定所有顶点都满足某条件,再切换执行另一个分支; - 选出特殊角色 —— 对顶点 ID 做
min或max,就选出了一个顶点来在算法里扮演特殊角色。
自定义方式是继承 Aggregator,指定聚合值如何从第一个输入值初始化、以及多个部分聚合值如何归约成一个;归约算子同样应当可交换且可结合。
一条容易忽略的差别:默认的 aggregator 只归约单个超步的输入值,但也可以定义 sticky aggregator,它使用所有超步的输入值 —— 例如维护一个只在增删边时才调整的全局边计数。
还能做更高级的事:用 aggregator 实现 min 聚合器,最小值在下一超步广播给所有 worker,最低桶号里的顶点去松弛边。
拓扑变更:偏序 + handler
有些图算法需要改图的拓扑:聚类算法可能把每个簇替换成一个顶点,最小生成树算法可能只保留树上的边。Compute() 除了发消息,也可以发出增删顶点或边的请求。
同一超步里多个顶点可能发出冲突的请求(例如两个请求以不同初始值添加同一个顶点
① 偏序。 变更和消息一样,在请求发出的下一个超步才生效;在那个超步内,顺序被固定下来:
先 删除
├─ 先删边
└─ 再删顶点 (删顶点会隐式删掉它的全部出边)
后 添加
├─ 先加顶点
└─ 再加边
最后 才轮到 Compute() 被调用这个偏序对大部分冲突都能产出确定性的结果。
② handler。 剩下解决不了的冲突交给用户定义的 handler。例如同一超步里多次请求创建同一个顶点时,默认行为是系统任意挑一个,有特殊需要的用户可以在 Vertex 子类里定义合适的 handler 来指定更好的冲突解决策略。多次删除顶点、多次增删边的冲突用的也是同一套 handler 机制。把解决方式委派给 handler 是为了让 Compute() 的代码保持简单 —— 代价是限制了 handler 与 Compute() 之间的交互,但实践中没有成为问题。
一条设计取向值得单独记:协调是惰性的 —— 全局变更在被应用之前不需要任何协调。这个选择有利于流处理。其直觉是:涉及顶点
另外,Pregel 支持纯局部变更,即一个顶点增删自己的出边、或者删除自己。这类变更不可能引入冲突,因此被做成立即生效,用更简单的顺序编程语义简化分布式编程。
输入与输出
图的文件格式有很多可能(文本文件、关系数据库里的一组顶点、Bigtable 里的行……)。为避免强加某种格式,Pregel 把"把输入文件解释成图"和"图计算"这两件事解耦;输出也能以任意格式生成、以最适合应用的形式存储。库里为常见格式提供了 reader 与 writer,特殊需求可以继承抽象基类 Reader / Writer 自己写。
输入的划分与加载
输入划分与图分区是两件正交的事,这个区分很实用,因为它解释了"为什么图看起来加载得很乱,却没有出错":
| 输入划分 | 图分区 | |
|---|---|---|
| 划分依据 | 文件的边界(一条条记录) | 顶点 ID( |
| 单位 | 一条记录,可含任意数量的顶点与边 | 一组顶点 + 它们的全部出边 |
| 谁来决定 | master 把输入分给各个 worker | master 定分区数,worker 各自持有若干分区 |
| 两者关系 | 互不依赖 —— 一条记录里的顶点可能散落在任意分区 |
于是加载过程要处理一种错位:从输入记录里读到的顶点,不一定归当前 worker 所有。两条路:
- 顶点归自己 → 立即更新本地数据结构,不产生任何通信;
- 顶点归别人 → 向拥有它的远端 peer 排一条消息 —— 走的是和计算期同一条消息通道,区别只在于这条"消息"承载的是一个顶点或一条边的初始状态。
加载完成后所有顶点被标记为 active,这是"超步 0 全 active"这条规则的来源。
一处值得记的工程后果:输入划分按文件边界这件事,决定了加载阶段的并行度上限。如果一个输入文件很大、而文件数少于 worker 数,就会有 worker 分不到输入、纯等;反过来文件碎成很多小份时,分片与调度开销又会上升。这是"输入格式解耦"这个设计选择的现实成本 —— 格式自由换来的是加载阶段的并行度需要用户自己保证。
实现
目标环境是 Google 的集群架构:成千上万台商用 PC 按机架组织、机架内带宽很高,机架之间互连但地理上分散。应用通常跑在一个集群管理系统上,由它调度作业以优化资源分配,有时会杀掉实例或把它们迁到别的机器;系统带名字服务,实例可以用逻辑名引用而不依赖它当前绑在哪台物理机上。持久数据放在分布式存储(GFS)或 Bigtable 上,临时数据(例如缓冲的消息)放在本地磁盘。
分区与执行阶段
分区的划分方式是:每个分区 = 一组顶点 + 这些顶点的全部出边。
一个性质很关键:顶点到分区的归属只取决于顶点 ID —— 这意味着即使某个顶点属于别的机器、甚至这个顶点还不存在,也能算出它该属哪个分区。默认的分区函数就是
顶点到 worker 机器的分配,是 Pregel 里分布不透明的唯一主要地方。 有些应用用默认分配就够,有些则能从自定义分配函数里获益,以更好利用图自身的局部性 —— 例如 Web 图上常见的启发式是把同一站点的页面放在一起。
无故障时,程序执行分几个阶段:
① 启动:多份用户程序副本起来,一个任 master(不持图,负责协调)
worker 通过名字服务发现 master,并发送注册消息
│
② 分配分区:master 决定分区数,给每个 worker 分一个或多个分区
每 worker 多于一个分区 → 分区间并行 + 更好负载均衡,通常提升性能
每个 worker 拿到「全部 worker 的完整分配表」
│
③ 分输入:master 把一部分输入分给每个 worker(输入是一条条记录,
每条含任意数量的顶点与边);输入划分与图分区正交,通常按文件边界
├─ 加载的顶点属于自己的分区 → 立即更新数据结构
└─ 否则 → 向拥有该顶点的远端 peer 排一条消息
加载完成后,所有顶点标记为 active
│
④ 循环超步:worker 遍历 active 顶点,每个分区一个线程,调用 Compute()
消息异步发送(重叠计算与通信、便于批量),但必须在超步结束前送达
结束后回报「下一超步将有多少 active 顶点」
重复直到没有 active 顶点、也没有在途消息
│
⑤ 停机后:master 可指示各 worker 保存自己那部分图容错:checkpoint + ping
容错通过 checkpoint 实现。 在每个超步开始时,master 指示 worker 把各自分区的状态存到持久存储 —— 包括顶点值、边值、以及入向消息;master 单独保存 aggregator 的值。
故障检测靠 master 定期向 worker 发的 ping 消息:如果 worker 在指定间隔内没有收到 ping,它自己终止;如果 master 收不到某个 worker 的回应,就把它标记为失败。
恢复流程:一个或多个 worker 失败时,它们所辖分区的当前状态就丢了。master 把图分区重新分配给当前可用的 worker 集合,它们从超步
checkpoint 的频率用一个"平均无故障时间"模型来选,在 checkpoint 成本与期望恢复成本之间做平衡。
受限恢复(confined recovery,开发中) 是针对恢复成本的改进:worker 在图加载与各超步期间额外记录自己分区的出向消息。于是恢复被限制在丢失的分区上(它们由 checkpoint 恢复),系统用健康 partition 的日志消息 + 恢复 partition 重算出来的消息,把缺失的超步重放到
受限恢复有一条硬前提:要求用户算法是确定性的 —— 否则"原执行留下的消息"与"恢复期新算出来的消息"混在一起会产生不一致。随机算法可以通过"基于超步与分区确定性地播种伪随机数发生器"来变成确定性的;真正不确定的算法可以关掉受限恢复,退回基本恢复机制。
checkpoint 频率怎么定,可以推一个简化的判据。设每个超步耗时为
- 摊销到每个超步的开销约为
; - 故障时期望重放的超步数约为
,重放代价是 。
前者随
越大(写全量分区状态越贵)→ 周期应该拉长; 越小(超步本身越快)→ 最优周期越短,因为重放的每一步都便宜; 没有出现在这个简化式子里 —— 它决定的主要是"值不值得做 checkpoint",而不是"多久做一次"。机器越可靠,"完全不 checkpoint"就越合理,这也解释了实测里关掉 checkpoint 的做法。
受限恢复多记了什么,值得对照着看:worker 在图加载期与每个超步期间额外记录自己分区的出向消息。故障发生后,系统手里有两份东西可用于重放 —— 健康分区留下的消息日志、以及被恢复分区重算出来的消息。于是恢复被限制在丢失的分区上,健康分区不用重算。
"磁盘带宽足以让 I/O 不成为瓶颈"这个前提有它的边界:出向消息的总量大致是"边数 × 超步数"量级,全部要顺序写到本地盘。它成立的前提是本地盘只承载这一类写入(持久数据在 GFS 或 Bigtable);同一块盘上还跑着别的 I/O 密集任务时,这个前提就不成立了。
一条排查线索:恢复时间偏长时,先分辨是"重放超步多"还是"重算分区多" —— 前者说明 checkpoint 太稀疏,后者说明失败波及的分区多、受限恢复没开启或不适用(算法非确定)。
三类需要分辨的故障形态,恢复路径各不相同:
| 故障 | 谁发现 | 走哪条路径 | 代价 |
|---|---|---|---|
| worker 进程挂了 | master(ping 收不到回应) | 分区重分配 + checkpoint 重载 + 重放 | 取决于 checkpoint 距失败点有多远 |
| 网络分区,两边互相听不见 | 双方都认为对方死了 | worker 自行终止,master 标记它失败 | 同上,且可能重复触发 |
| 集群管理系统主动杀掉或迁移实例 | 同上 | 同上 | 同上 |
第二、三行是同一套实现要同时应付的场景 —— 运行环境本身有时会杀掉实例、或把它们迁到别的机器。这就是"故障是常态而不是例外"这个假设的具体来源,容错机制存在的理由落在这里。
恢复完成后值得确认三件事:用的 checkpoint 是"最近可用"的那个("最近可用"这几个字意味着更近的 checkpoint 可能本身就不完整);重放确实从超步
worker 内部
worker 把图的一部分状态放在内存里:可以看成一个从顶点 ID 到顶点状态的映射,每个顶点的状态由四样东西组成 —— 当前值、出边表(目标顶点 ID + 边值)、入向消息队列、一个标明是否 active 的标志。
执行一个超步时,worker 遍历所有顶点调用 Compute(),传入当前值、消息迭代器、出边迭代器。入边无法访问 —— 每条入边是源顶点所持有的列表的一部分,而那个源顶点通常在另一台机器上。
有三处实现细节值得记:
① active 标志与消息队列分开存(出于性能考虑)。
② 顶点和边的值只有一份,但 active 标志与入向消息队列有两份 —— 当前超步一份、下一超步一份。原因是并发性:worker 在处理超步
③ 发送路径分远端与本地:
- 远端 —— 消息先缓冲;当缓冲大小达到阈值时,最大的那些缓冲被异步刷出,每个作为一条网络消息投递到目标 worker;
- 本地 —— 可以直接放进目标顶点的入向消息队列(一个优化)。
若用户提供了 Combiner,它在两个位置都会被应用:消息加入出向队列时,以及到达入向队列时。后一个位置不减少网络用量,但减少存储消息所需的空间。
master 内部
每个 worker 在注册时得到一个唯一标识符。master 维护一张当前存活 worker 的列表,含唯一标识符、寻址信息、以及它被分配到图的哪一部分。一个重要的规模性质:master 数据结构的大小与分区数成正比,而与顶点数、边数无关 —— 所以单个 master 就能协调哪怕非常大的图。
大多数 master 操作在屏障(barrier)处终止:向"操作开始时已知存活"的每个 worker 发出同一个请求,并等待每个 worker 的响应。任一 worker 失败就进入恢复模式;屏障同步成功则进入下一阶段 —— 对计算屏障而言,就是递增全局超步号再进入下一个超步。
master 还维护一批运行统计:图的总大小、出度分布直方图、active 顶点数、近期超步的耗时与消息流量、以及所有用户定义 aggregator 的值。为了让用户能监控,master 跑一个 HTTP 服务器展示这些信息。
aggregator 的实现
每个 worker 维护一组 aggregator 实例,用类型名加实例名标识。一个 worker 为图的任一分区执行超步时,先把供给某 aggregator 实例的所有值合并成一个局部值 —— 也就是"在该 worker 该分区的全部顶点上部分归约"的结果。超步结束时,worker 之间组成一棵树,把部分归约的 aggregator 归约成全局值并交给 master。
这里有个刻意的选择:用树形归约,而不是沿 worker 串成一条链做流水线 —— 目的是在归约过程中并行使用 CPU。最后 master 在下一个超步开始时把全局值发给所有 worker。
一个超步里的时序:从屏障到屏障
"所有顶点并行执行同一个函数"是一句抽象描述。把它摊到时间轴上,一个超步是一串有序动作,而这些动作的次序解释了模型里几乎所有性能性质的来源:
① master 递增全局超步号,把「执行下一超步」下发给已知存活的每个 worker
│ (带上上一超步归约出的 aggregator 全局值)
▼
② worker 以「每个分区一个线程」遍历 active 顶点,逐个调 Compute()
├─ 读上一超步到达的消息(迭代器)
├─ 改写顶点值与出边值
├─ SendMessageTo():本地 → 直投入向队列;远端 → 进缓冲
└─ 可能 VoteToHalt(),也可能请求增删点边
▼
③ worker 收尾:把出向缓冲刷出;本超步该到达的消息必须全部落进入向队列
▼
④ worker 回报 master:「下一超步将有多少 active 顶点」
▼
⑤ master 收齐全部回报 → 屏障成立 → 回到 ①超步耗时由最慢的 worker 决定
消息"异步发送"与"必须在超步结束前送达"是一对约束。异步指的是 SendMessageTo() 立刻返回、不阻塞计算线程,发送与计算可以重叠;但③那一句把它拉了回来 —— 屏障不只等计算结束,还要等通信结束。于是"超步耗时"的准确写法是:
这个 max 是理解顶点中心性能的钥匙:任何"平均"指标都掩盖了真正的决定量。它同时解释了为什么"计算与通信重叠"的收益是有限的 —— 发送可以重叠、接收不能,因为下一个超步的 Compute() 一执行就要把消息迭代器交出去,递不出来就没法开始。
回报值同时承担两个职责
④ 的回报值既是下一个超步的工作量估计(master 据此判断要不要继续),也是终止条件的判据 —— 全部 worker 都回报"0 个 active 顶点"、且没有消息在途,算法才整体停机。
这也回过头解释了投票停机的第三条规则:顶点被消息唤醒后必须显式再停用自己,否则下一次回报的计数不会归零,算法就停不下来。规则看似琐碎,落点却在这里。
底层依赖:消息的序列化、缓冲与落盘
"网络只用于传消息"这句话落到物理层,要经过三个落点,每一处的容量都有限:
计算线程 SendMessageTo()
│
① 出向缓冲(内存,按目标 worker 分组)
│ 缓冲大小达到阈值 → 挑最大的几个异步刷出,各作为一条网络消息
▼
② 网络传输(机架内带宽高,机架间与跨集群受地理分布限制)
▼
③ 目标 worker 的入向消息队列(内存,按顶点分组)① 缓冲为什么按"目标 worker"分组、又按阈值刷出:一条消息单独发一次的固定开销远大于消息本身,攒起来发才能摊薄;而"挑最大的那些先发"是为了让大缓冲尽早进入网络、不至于全部卡在最后。临时数据(缓冲中的消息)落在本地磁盘,持久数据(输入图、输出、checkpoint)落在 GFS 或 Bigtable —— 这个分工意味着本地盘的延迟不进入关键路径,因为它只承载暂存。
② 本地消息有一条直投优化:目标顶点归本 worker 所有时,消息直接进目标顶点的入向队列,不经过缓冲与网络。这条优化顺带提示了"合并该在何处生效" —— 若用户提供了 Combiner,它在两个位置都会被应用:消息加入出向队列时、以及到达入向队列时。后一个位置不减少网络流量,但减少缓冲占用的空间。想省网络要在第一个位置生效,想省内存则两个位置都起作用。
③ "每顶点一个值"这个限制,出口在序列化格式上。顶点值、边值、消息值都是模板参数,看起来一板一眼;实际写法是用协议缓冲区这类能承载任意结构的类型把它们顶上去 —— 于是"每顶点只有一个值"约束的是值的数量,不约束值的结构复杂度。代价是这些值都要经过序列化才能跨机器、才能落盘,而序列化开销是顶点中心框架里一处固定的成本项。
拓扑假设决定优化方向:集群由成千上万台商用 PC 按机架组织,机架内带宽很高,机架间互连但地理上分散。于是"顶点放在哪台机器"直接决定跨机架流量,而顶点到 worker 的分配是整个模型里唯一不透明的地方。默认的哈希分配在多数图上够用;图自身带局部性时(Web 图上把同一站点的页面放在一起)自定义分配能拿到实打实的收益。
内存是这套设计的硬边界:worker 把图的这一部分整个放在内存里 —— 概念上是"顶点 ID → 顶点状态"的映射,每个顶点状态含当前值、出边表、入向消息队列、active 标志四样。入边无法访问的物理原因也落在这里:入边是源顶点所持列表的一部分,而那个源顶点通常在另一台机器上。
这也解释了 worker 内部两处实现细节为什么不是随意的:顶点值与边值只有一份,active 标志与入向消息队列却各存两份(当前超步一份、下一超步一份)。原因是并发 —— worker 在处理超步
把这条链读完,一个结论是自然浮出来的:整个计算状态都在 RAM 里,"图装不装得下"直接决定作业能不能跑 —— 这正是后来 out-of-core 与单机磁盘路线要解决的事。
序列化的成本出现在三处,加起来才是"顶点中心比单机实现慢"的那部分固定开销:
| 出现的位置 | 为什么必须序列化 | 能不能避免 |
|---|---|---|
| 跨机器的消息 | 要经过网络 | 不能,只能靠 Combiner 减少条数 |
| 跨超步的持久化(checkpoint) | 要写进持久存储 | 可以靠降低 checkpoint 频率摊薄 |
| 图加载与输出 | 输入输出走 GFS 或 Bigtable | 不能,但只发生一次 |
内存占用可以粗估,估法值得记下来,因为它决定分区数:把"顶点数 × 每顶点状态开销 + 边数 × 每边开销"加起来,再乘一个系数容纳同一时刻的两份 active 标志与两份消息队列。这个估算的用处很直接 —— 分区数不是越大越好:分区多则每 worker 装得少、负载均衡好,但每个分区都有一份固定开销(消息队列、迭代器、线程栈),分区太碎时这份开销会盖过收益。
一条实际的观察路径:内存吃紧时先看"边数占了多少"。出边表是每个顶点状态里最容易膨胀的一块,而它的大小与出度成正比 —— 出度分布呈重尾时,少数高度顶点的出边表可能吃掉可观的一份内存。
参数与可调项
Pregel 本身是 Google 内部系统,没有公开的配置文档页,所以这一节分两栏看:模型层定下了哪些旋钮,以及这些旋钮在开源对应物(Apache Giraph)里被具体化成什么参数名与默认值。下面的默认值照录 Giraph 官方选项页,不做推测。
模型层定义的旋钮有七个,归属各不相同:
| 旋钮 | 影响什么 | 由谁决定 |
|---|---|---|
| 分区函数 | 顶点落到哪个分区(默认 | 用户可替换 |
| 分区数 | 并行度、每个 worker 拿几个分区 | 框架默认,可覆盖 |
| 顶点到 worker 的分配 | 跨机架流量 | 用户可替换(唯一不透明处) |
| Combiner | 消息条数与缓冲占用 | 用户声明,默认关闭 |
| Aggregator | 全局归约的语义与生命周期(是否 sticky) | 用户声明 |
| 冲突 handler | 拓扑变更里偏序解决不了的冲突 | 用户声明,默认任选一个 |
| checkpoint 频率 | 恢复成本与运行时开销的平衡 | 由"平均无故障时间"模型定 |
在 Giraph 里,这些旋钮长成下面这样:
| 参数 | 默认值 | 语义 |
|---|---|---|
giraph.numComputeThreads | 1 | 顶点计算用几个线程 |
giraph.numInputThreads / giraph.numOutputThreads | 1 / 1 | 输入切分加载、输出写出的线程数 |
giraph.graphPartitionerFactoryClass | HashPartitionerFactory | 顶点到分区的划分(对应上面的分区函数) |
giraph.partitionClass | SimplePartition | 分区容器实现 |
giraph.userPartitionCount | -1 | -1 表示按框架默认算,非负则覆盖分区数 |
giraph.useOutOfCoreGraph | false | 是否把图放到内存之外 |
giraph.outEdgesClass | ByteArrayEdges | 出边的存储实现,决定边值怎么落内存 |
giraph.messageEncodeAndStoreType | BYTEARRAY_PER_PARTITION | 消息的编码与存储方式 |
giraph.checkpointFrequency | 0 | 0 表示不 checkpoint,1 表示每个超步都做,2 表示每两个超步 |
giraph.checkpointDirectory | _bsp/_checkpoints/ | checkpoint 落在 HDFS 的哪个目录 |
giraph.maxNumberOfSupersteps | 1 | 完成这么多个超步后停机 |
giraph.SplitMasterWorker | true | master 与 worker 分离,动态恢复需要它 |
giraph.masterComputeClass | DefaultMasterCompute | master 侧的计算逻辑 |
giraph.metrics.enable | false | 是否开启 Metrics |
giraph.logLevel | info | 覆盖 Hadoop 日志级别 |
按六要素把其中三个最要紧的摊开。
giraph.checkpointFrequency(默认 0) —— 语义是"每多少个超步做一次 checkpoint"。默认值本身就是一个立场:0 意味着开箱状态下不落 checkpoint,失败的代价是整个作业从头重跑。什么时候该改?图大到"重算一次的成本不可接受"时。往小改(更频繁)→ 恢复更快,但每个超步都要把分区状态全量写进持久存储,运行时开销随之上升;往大改(更稀疏)→ 运行时开销降低,但单次恢复要重放的超步更多。联动的是 checkpointDirectory 所在存储的写入带宽。失败模式两头都像"没配过这个参数":配得太密时表现为超步耗时莫名变长(时间花在写盘上),配得太疏时表现为机器一挂就要重跑很久。
这里有一处值得对照的地方:模型层面把 checkpoint 定在"每个超步开始时",而开源实现的默认是关掉。同一件事在两处的默认值相反,说明"机制存在"与"默认启用"是两个独立的决定 —— 而且 Pregel 自己的实测就是关掉 checkpoint 跑的(失败概率低、跑得快)。
giraph.maxNumberOfSupersteps(默认 1) —— 语义是"完成这么多个超步后停机"。官方描述给了一个例子:设成 3,则作业最多跑超步 0、1、2,然后进入 shutdown 超步。按这个语义,默认 1 意味着不显式改它就只能跑超步 0 就收尾 —— 也就是"加载图 + 一轮计算"。需要多轮迭代的算法(PageRank、最短路、半聚类)都必须显式把这个值放大。失败模式很直白:作业秒退、超步计数停在 1,而日志层面不报错,看起来像"图没加载进来"。
giraph.useOutOfCoreGraph(默认 false) —— 语义是"把图放到内存之外"。默认 false 对应模型原本的假设:整个计算状态常驻 RAM。什么时候该改?图大到单机内存装不下、又拿不到 tb 级内存时。改成 true 的后果是计算期间反复读写磁盘,超步耗时上一个台阶;不改的后果是装不下就起不来。这条参数之所以存在,直接对应模型自己承认的一处短板:"目前整个计算状态都在 RAM 里,已经有一部分数据 spill 到本地盘,会继续往这个方向走。"
几条联动关系值得记住:
numComputeThreads与partitionClass:模型层的说法是"每个分区一个线程",实现层把它显式化成线程数 —— 分区数与线程数的配比决定 CPU 用不用得满;- 分区数与
graphPartitionerFactoryClass:换分区函数往往要连带调分区数,否则"顶点分布变了、每分区大小也变了"; - Combiner 与
messageEncodeAndStoreType:Combiner 减少消息条数,编码方式影响每条消息的字节数,两者作用在不同维度上,可以叠加。
三组容易混淆的机制
Pregel 的机制里有三对概念,名字上都带"合并"或"解决冲突"的意思,但作用域差得远。摆成表比分散在各节里读要清楚。
Combiner 与 Aggregator
| Combiner | Aggregator | |
|---|---|---|
| 合并的对象 | 发往同一个目标顶点的多条消息 | 全图每个顶点各贡献的一个值 |
| 谁看得到结果 | 只有那个目标顶点 | 所有顶点 |
| 生效位置 | 发送路径上两处(出向队列、入向队列) | 超步结束时,worker 之间树形归约 |
| 生效时机 | 同一个超步内(消息还没被读走) | 下一个超步 |
| 默认是否开启 | 否 | 预定义的 min/max/sum 直接可用 |
| 成立前提 | 操作可交换、可结合 | 归约算子可交换、可结合 |
一句话概括分工:Combiner 压的是消息的条数,Aggregator 换的是全局的视角。
消息与 Aggregator
| 消息 | Aggregator | |
|---|---|---|
| 拓扑 | 点对点,目标可以完全任意 | 全收集 → 归约 → 全广播 |
| 目标能否是非邻居 | 能 | 不适用,没有"目标"这个概念 |
| 收到的内容 | 发送方决定的具体值 | 一个归约结果 |
| 延迟 | 一个超步 | 一个超步 |
这条对照的实用价值在于选择:需要"全图某个量的汇总"时,用消息模拟要每个顶点发给同一个收集者再广播回去 —— 两轮,且收集者会成瓶颈;用 Aggregator 是一轮归约加一轮广播,归约本身还是树形并行的。
偏序与 handler
| 偏序 | handler | |
|---|---|---|
| 覆盖的冲突 | 大部分 | 偏序解决不了的剩余部分 |
| 由谁定义 | 框架内建 | 用户定义 |
| 典型场景 | 同一超步里既删又增、删顶点隐含删它的出边 | 多次请求创建同一个顶点 |
| 默认行为 | —— | 任意挑一个 |
先偏序、后 handler,是两层依次兜底的关系:能靠时序定下来的就不打扰用户,定不下来的才把决定权交出去。把冲突解决委派给 handler 的代价是限制了 handler 与 Compute() 之间的交互,换来的是 Compute() 的代码保持简单。
还有一处容易混:sticky 与普通 Aggregator。普通 Aggregator 只归约本超步的输入值,sticky 的累积所有超步的输入值(例如维护一个只在增删边时才调整的全局边计数)。选错的表现是"全局计数每轮被重置",而这类错误在结果上往往只体现为一个偏小的数字。
四个应用
PageRank
顶点值类型 double(存暂定的 PageRank),消息类型 double(承载 PageRank 分量),边值类型是 void —— 因为边不存信息。
class PageRankVertex : public Vertex<double, void, double> {
public:
virtual void Compute(MessageIterator* msgs) {
if (superstep() >= 1) {
double sum = 0;
for (; !msgs->Done(); msgs->Next()) sum += msgs->Value();
*MutableValue() = 0.15 / NumVertices() + 0.85 * sum;
}
if (superstep() < 30) {
const int64 n = GetOutEdgeIterator().size();
SendMessageToAllNeighbors(GetValue() / n);
} else {
VoteToHalt();
}
}
};初始状态假定超步 0 里每个顶点的值是
到第 30 个超步之后不再发消息,每个顶点投票停机。现实中 PageRank 会一直跑到收敛,而"检测收敛条件"正是 aggregator 的用武之地。
单源最短路
最短路有三个重要变体:单源(一个源到其余每个顶点)、s-t(给定
s-t 这类问题实际上相对容易:典型路网上的解法只访问极小一部分顶点 —— 有一例观察到的是3200 万顶点里只访问了 8 万个。这里选单源变体讲,因为它更贴合大规模图,同时又能给出更有意思的扩展性数据。
class ShortestPathVertex : public Vertex<int, int, int> {
void Compute(MessageIterator* msgs) {
int mindist = IsSource(vertex_id()) ? 0 : INF;
for (; !msgs->Done(); msgs->Next())
mindist = min(mindist, msgs->Value());
if (mindist < GetValue()) {
*MutableValue() = mindist;
OutEdgeIterator iter = GetOutEdgeIterator();
for (; !iter.Done(); iter.Next())
SendMessageTo(iter.Target(), mindist + iter.GetValue());
}
VoteToHalt();
}
};每个顶点的值初始化为 INF(一个大于从源点出发任何可行距离的常数)。每个超步里,顶点先收到邻居发来的"候选最小距离";若这些更新里的最小值小于当前值,就更新自身并向邻居发出"每条出边的权重加上新最小值"。**第一个超步里只有源顶点会更新(从 INF 变成 0)并向它的直接邻居发出更新,这些邻居再更新并向下传播 —— 于是形成一道波前(wavefront)**在图中推进。不再有更新时终止,此时每个顶点的值就是从源点到它的最小距离(INF 表示不可达)。所有边权非负时保证终止。
这个算法的消息本质是"可能更短的距离",而接收方最终只关心其中的最小值,所以它适合用 Combiner 优化:
class MinIntCombiner : public Combiner<int> {
virtual void Combine(MessageIterator* msgs) {
int mindist = INF;
for (; !msgs->Done(); msgs->Next())
mindist = min(mindist, msgs->Value());
Output("combined_source", mindist);
}
};这个 Combiner 大幅减少了 worker 之间传输的数据量,也减少了执行下一个超步之前缓冲的数据量。
一条诚实的取舍:这个算法会比顺序版本(Dijkstra、Bellman-Ford)做多得多的比较,但它能在任何单机实现都无法企及的规模上求解最短路。更高级的并行算法(如 Thorup、
"波前"这个词值得坐实,因为它决定了这个算法在什么图上快、什么图上慢。波前是"本轮被更新的顶点集合",它的规模取决于图的结构:直径小、分支多的图(社交网络)上波前迅速铺满全图,前几个超步就把大部分顶点更新完;直径大的链状或网格状图上波前始终很窄,更新的轮数接近图的直径。这解释了一条实测规律:在平均出度低的图上,运行时间随图规模近似线性增长 —— 二叉树正是"波前窄、轮数随深度增长"的形态。
与顺序实现的关系可以用一张表摆清:
| Dijkstra | Bellman-Ford | Pregel 版 | |
|---|---|---|---|
| 更新依据 | 从已定顶点贪心扩展 | 逐轮松弛全部边 | 逐超步松弛"被本轮更新过"的顶点的出边 |
| 比较次数 | 每条边摊到 | 多于 Dijkstra,少于 Bellman-Ford 的全量松弛 | |
| 终止 | 目标确定即停 | 固定 | 没有更新就停 |
| 规模上限 | 单机内存 | 单机内存 | 单机装不下的图 |
一条边界必须写清:所有边权非负才保证终止。边权为负时,"更短的路径被反复发现"会让顶点反复唤醒,环上的负权尤其危险 —— 顶点可以在两个超步之间来回更新而收敛不了。这是"没有更新就停"这个终止条件在负权下的失效,落到实现上是找不到可修的缺陷的。
消息侧的量化:接收方最终只关心其中的最小值,所以 MinIntCombiner 能把发往同一顶点的多条候选距离压成一条。收益随"一个顶点平均被更新几次"放大 —— 更新次数越多,同一目标顶点在同一个超步里收到的候选越多,能合并掉的就越多。
二分图匹配
输入由两组互不相交的顶点构成,边只跨组;输出是一个端点互不相同的边子集。极大匹配是指无法在不共享端点的前提下再加入任何一条边的匹配。这里讲其中的随机化极大匹配算法。
顶点值是一个二元组:标明顶点属于 void,消息是布尔值。
算法的骨架是四阶段一轮的循环,阶段号就是"超步号 mod 4",整体是一个三方握手(WL 表示左组顶点):
各阶段的细节值得逐个说清,因为停机的时机是这套四阶段设计的核心:
- 阶段 0:每个未匹配的左顶点向它的每个邻居发一条匹配请求,然后无条件投票停机。这里有一个微妙的判断:若它没有发出任何消息(已经匹配了,或者没有出边),或者消息的所有接收者都已经匹配,那么它再也不会被重新激活;否则它会在两个超步后收到回应而被激活。
- 阶段 1:每个未匹配的右顶点随机选择收到的一条消息,给它批准,并给其余请求者发拒绝;然后无条件投票停机。
- 阶段 2:每个未匹配的左顶点从收到的批准里选一条,发出接受消息。已经匹配的左顶点永远不会执行到这个阶段 —— 因为它们在阶段 0 就没发过消息。
- 阶段 3:一个未匹配的右顶点最多收到一条接受消息。它记下匹配对象,然后无条件投票停机 —— 它没有别的事要做了。
半聚类
Pregel 被用于多个版本的聚类。其中半聚类(semi-clustering)出现在社交图上:顶点代表人,边代表人与人之间的连接,连接可能基于显式动作(例如在社交网站上加好友),也可能从行为里推断(例如邮件往来、合作发表);边可以带权,表示互动的频率或强度。
半簇是"一群互相频繁互动、与外部互动较少的人"。它与普通聚类的区别在于:一个顶点可以属于多个半簇。
算法是并行的贪心:
- 超步 0:
把自己作为规模 1、得分 1 的半簇写进自己的列表,并把自身发布给所有邻居; - 后续超步:
遍历上一超步收到的半簇 。若某个半簇 还不包含 、且 ,就把 加入它形成 ; - 把
与 按得分排序,把最好的那些发给 的邻居; 用"包含 的那些半簇"更新自己的列表。
终止条件是半簇不再变化,或者(为了性能)达到用户指定的超步上限。此时可以把每个顶点的最佳半簇候选聚合成一个全局的最佳半簇列表。
实测:单源最短路
环境是 300 台多核商用 PC,所有边权隐含为 1。不计集群初始化、在内存中生成测试图、以及校验结果的时间。因为所有实验都能在较短时间内跑完、失败概率低,checkpoint 被关掉了。
先看按 worker 数扩展:十亿顶点的二叉树(因此是十亿减一条边),worker 数从 50 变到 800 —— 运行时间从 174 秒降到 17.3 秒,16 倍的 worker 换来约 10 倍加速。
再看按图规模扩展:二叉树从十亿到五百亿顶点,固定 800 个 worker —— 时间从 17.3 秒增到 702 秒,说明在平均出度低的图上,运行时间随图规模线性增长。
二叉树显然不代表真实图,所以另外用对数正态分布的出度做实验:
取
这些数字不该被读成 Pregel 的最好成绩,有两个原因:实验中用的是默认的随机 hash 分区函数(拓扑感知的分区函数会给出更好性能),而且用的是朴素的并行最短路算法(更高级的算法会更好)。想说明的是:用相对很少的编码量就能拿到令人满意的性能。 十亿顶点与边的结果与 Parallel BGL 在 112 个处理器 / 2.56 亿顶点 / 10 亿边上的
版本与后续演进
Pregel 自己没有后续版本 —— 它是在生产里用了几年、把模型和实现讲清楚就停在那里的一个系统。真正在演进的是它留下的模型,而设计里主动承认了三处不足,此后十年的顶点中心系统基本可以按这三条线归类。
不足一:同步屏障的代价。 每个超步结束时所有 worker 都要等齐,快的那批必须空等慢的。沿这条线的做法是把同步放松:
- GraphLab(2012)改用异步执行,顶点主动拉取邻居数据而不是被动收消息;
- PowerGraph(2012)进一步按边切分而不是按顶点切分,缓解高度顶点造成的负载倾斜;
- GiraphUC(2015)走折中,把同步与异步混起来。
这里有一处值得记的反直觉结果:异步模型对随机游走这类任务收敛更快,但在多数任务上比同步慢,主要开销落在加锁与解锁上 —— 放松同步换来的并行度,可能还不够抵消一致性维护的成本。
不足二:整个计算状态常驻 RAM。 这条线的第一步是开源实现的 out-of-core 模式;更彻底的一条是把磁盘当作一等存储、用单机处理巨图(GraphChi、X-Stream、VENUS 这一类)。单机路线的代价也很明确:每轮迭代都要扫一遍磁盘上的图,哪怕本轮只有很小的子集需要计算;对"重迭代、轻载"的分析型作业尚可,对轻量查询就很浪费。
不足三:顶点到机器的分配。 这是模型自己点名的难点 —— 按拓扑分区在拓扑与消息流量一致时够用,不一致时就不够。三个系统各补了一块:
- GPS(2013):把高度顶点镜像到多台机器上,让"给自己全部邻居发同一条消息"这个动作不必跨机器摊开;
- Mizan(2013):运行时监控 + 顶点迁移做动态负载均衡,处理的是"分区时算不准、运行时才暴露"的倾斜;
- Pregel+(2015):把镜像与消息合并放在一起,用代价模型判断某个顶点值不值得镜像,另加一套请求-响应机制,专门压掉"只读远端值"这一类开销。
另一条方向是把顶点中心收进通用抽象,这也是 RDD 那篇预告过的事:
- GraphX(2014)在 Spark 上提供名为
pregel的算子 —— 但它实现的是 GAS(gather-apply-scatter),而非逐顶点的消息推送; - Pregelix(2015)把顶点中心架在数据流引擎上;Gelly(2015)在 Flink 上同时给出顶点中心、scatter-gather 与 GAS 三套 API。
还有一支走子图中心 / 块中心:Giraph++(2013)把子图整体暴露给用户编程,Blogel(2014)让块可以像顶点一样有状态并通信 —— 它们针对的是"顶点中心把消息条数放大到边数量级"这个痛点,用"把一组顶点当成一个单元"来压通信量。
作为这条线里唯一还在持续发版的开源实现,Apache Giraph 的发布节奏可以当作顶点中心系统的生命周期参照:0.1-incubating(2012-02)→ 1.0.0(2013-05)→ 1.1.0(2014-11)→ 1.2.0(2016-10)→ 1.3.0(2020-06)。它在基本模型之外补的东西集中在四处:master 侧计算、分片 aggregator、面向边的输入、out-of-core。
一条选型判据:先问"瓶颈落在哪一条线上",答案直接指向该看哪类系统 —— 瓶颈是屏障空等就看异步或混合路线;是内存装不下就看 out-of-core 或单机磁盘路线;是少数高度顶点拖慢全局就看镜像与迁移这一类优化。而如果作业本身不是重迭代的,通用引擎上的一个 pregel 算子往往够用,不必为此另起一套专用集群。
回看这三条线,有一个共同点值得记:它们都在改动"顶点中心"这个模型的外围,而没有动它的内核。内核始终是"顶点持状态、超步同步、消息通信"这三件事;被改的是同步的松紧(异步与混合)、状态的落点(内存与磁盘)、以及顶点的分布(镜像、迁移、块)。
这个观察对选型有用:要判断一个图系统适不适合手上的作业,看它的内核就够了 —— 顶点能不能改图拓扑、消息能不能发给非邻居、超步是否同步。外围优化只影响常数因子,内核决定哪些算法写得出来。按这个口径,只有实现了完整模型语义(含拓扑变更、任意目标的消息、超步同步)的系统才能直接承接一份顶点中心程序;不支持图变更的那批系统(GraphLab、GraphX 的 pregel 算子等)遇到"算法中途改图"这一类问题就不适用。
另一条演进方向不在模型上,而在抽象层级上:子图中心与块中心把"一组顶点"作为编程单元,针对的是消息条数被放大到边数量级这个结构性问题。它们的取舍也很清楚 —— 换到更高的抽象粒度可以压通信量,代价是把"哪些顶点需要一起算"这个决定提前到分区阶段,不再像顶点中心那样交给运行时。
与 BSP 的关系:借了什么、多做了什么
高层组织方式借鉴 BSP(Bulk Synchronous Parallel)。BSP 把一个并行计算拆成超步,每个超步含三件事 —— 本地计算、通信、一次全局同步;模型的参数是三个量:处理器数、每轮通信的能力、同步的间隔。Pregel 的超步模型直接来自这里,包括"超步之间由全局同步点分隔"这一条。
| BSP 的抽象 | Pregel 在它之上加的东西 | |
|---|---|---|
| 计算单元 | 通用的"处理器" | 顶点状态机(顶点值 + 出边 + 消息队列 + active 标志) |
| 通信原语 | 由库提供的一组原语 | 点对点消息,目标可任意指定 |
| 全局协调 | 没有内建 | Aggregator |
| 终止判定 | 用户自己判断轮数 | 投票停机 + 无消息在途 |
| 容错 | 不在模型里 | checkpoint + ping,以及受限恢复 |
| API 层级 | 通信原语 | 图专用 API(Compute() 一个函数) |
差别最大的一行:容错
已有的通用 BSP 库实现(Oxford BSP Library、Green BSP、BSPlib、Paderborn BSP)在通信原语、可靠性处理、负载均衡与同步上的取舍各不相同,但它们的可扩展性与容错没有被评估到几十台机器以上,而且没有一个提供图专用 API。Pregel 的贡献正落在这一格:把 BSP 的超步骨架套在顶点状态机上,再把容错做进系统。
一条容易忽略的继承关系:BSP 的"全局同步"在 Pregel 里变成了"消息可见性的边界"。BSP 里同步是为了让各处理器的通信结果对下一轮可见;Pregel 里超步屏障起的是同一个作用 —— 超步
与相关工作的边界
与 MapReduce 概念上相似,但 Pregel 提供自然的图 API,且对图上的迭代计算支持得好得多。这一点也让它区别于其他隐藏分发细节的框架(Sawzall、Pig Latin、Dryad)。另一处更根本的差别是模型:
Pregel 实现的是有状态的模型 —— 长生命周期的进程做计算、通信并修改本地状态;而不是数据流模型 —— 任何进程只对输入数据计算、产出被别的进程消费的输出。
它受 BSP 启发(超步模型就来自这里)。已有的通用 BSP 库实现(Oxford BSP Library、Green BSP、BSPlib、Paderborn BSP)在提供的通信原语、以及如何处理可靠性(机器故障)、负载均衡与同步上各不相同;但它们的可扩展性与容错没有被评估到几十台机器以上,而且没有一个提供图专用的 API。
最接近的两个:
| 系统 | 差别 |
|---|---|
| Parallel BGL | 用 property map 存顶点与边的信息,用 ghost cell 保存远端组件的值 —— 需要引用大量远端组件时会出现扩展问题;Pregel 用显式消息获取远端信息,不在本地复制远端值。最关键的区别是Pregel 提供容错,因此能在"故障很常见(硬件故障、或被更高优先级作业抢占)"的大集群环境里工作 |
| CGMgraph | 用 CGM(Coarse Grained Multicomputer)模型 + MPI;底层分发机制对用户暴露得多得多;侧重"给出算法实现"而不是"提供一个用来实现算法的基础设施";用面向对象风格而非泛型风格,有一定性能代价 |
除 Pregel 与 Parallel BGL 之外,很少有大到数十亿顶点量级的实验结果,而其中最大的那些来自定制的 s-t 最短路实现而不是通用框架:BlueGene/L 上 32,768 个 PowerPC 处理器跑广度优先搜索,3.2 十亿顶点 / 320 亿边的 Poisson 随机图用 1.5 秒;Cray MTA-2 上 10 节点的高多线程系统,1.34 亿顶点 / 8.05 亿边的 R-MAT 图用 0.43 秒;Parallel BGL 在 200 处理器的 x86-64 Opteron 集群上,40 亿顶点 / 200 亿边的 Erdős–Rényi 随机图用 0.43 秒 —— 他们把更好的性能归因于 ghost cell,同时观察到自己的实现在超过 32 个处理器后性能开始变差。
停机条件:写不对会怎样
终止判定看起来只有三条规则,但写错的形态集中且可枚举。这是顶点中心程序里唯一一处"写成什么样都不报错、只会跑得不对"的地方,值得单独摆开:
| 写法 | 后果 | 表现 |
|---|---|---|
该 VoteToHalt() 时没调 | 顶点永远 active,算法不终止 | 超步数一路涨到上限;没设上限则跑到资源耗尽 |
不该 VoteToHalt() 时调了 | 该被唤醒的顶点提前退出 | 结果偏小或部分顶点没被更新,且不报错 |
| 在发送消息之前停用自己 | 消息照发、顶点照停,唤醒语义被搞混 | 结果对,但轮数比预期多(每次都多一次纯唤醒) |
| 只按"没有消息"判终止 | 忽略"消息在途"这个条件 | 偶发提前终止,结果不稳定(同一份输入两次结果不同) |
第三行值得展开:发送消息与投票停机是两个独立动作,次序不影响消息是否发出(消息在发送时刻就进了队列),但影响顶点何时被重新唤醒。稳妥的写法是在超步末尾统一决定是否停用,别在中途根据某个分支条件提前停。
最后一条规则的后果
第四行是那条容易漏掉的规则的具体后果 —— 算法整体终止的条件是"所有顶点 inactive 且没有消息在途",两个条件缺一不可。只判前一半的写法在单机上往往看不出问题,在真集群上会因为消息延迟而偶发提前终止。
一条通用判据:停机条件的正确性可以靠"多跑一遍"来验 —— 同一份输入与同一个算法跑两次,结果不一致通常指向终止判定,而非计算逻辑。
排查:从症状到判据
顶点中心作业的运行特征和批处理很不一样:超步屏障把"最慢的 worker"变成了超步耗时本身,所以观察的落点落在每个超步的耗时曲线与 active 顶点数曲线上,而不在吞吐。
| 症状 | 判据 | 先看什么 | 常见归因 |
|---|---|---|---|
| 超步耗时呈阶梯状(稳定若干轮后突然上一个台阶) | 某几个 worker 的顶点数或消息量远高于其他 | 分区分布、出度直方图 | 倾斜 —— 少数高度顶点挤在同一个分区 |
| 超步耗时不降、active 顶点数也不降 | 每轮都有大量顶点在给别人发消息 | active 计数、消息总量 | 停机条件写错,或收敛判据没落到 VoteToHalt() 上 |
| 消息量远大于顶点数 | 消息条数与顶点数的比值 | 消息流量统计 | 消息没聚合,Combiner 未启用 |
| 部分 worker 早已空闲,超步却不结束 | 屏障在等最慢的那个 | 各 worker 的完成时刻 | 负载不均,或某台机器本身慢 |
作业秒退,superstep 停在 1 | 是否显式设过超步上限 | 启动参数 | 超步上限取了默认值 |
| 图加载完就 OOM | 单机要装的顶点数与出边表规模 | 顶点数、边数、分区数 | 状态常驻内存;需减少每机分区或改 out-of-core |
| 机器挂掉后重跑很久 | 上一个可用 checkpoint 距失败点有多远 | checkpoint 目录、checkpoint 频率 | checkpoint 关了或太稀疏 |
| 恢复期结果与原执行不一致 | 算法是否确定 | 受限恢复是否开启 | 算法里有非确定性(随机、读外部状态) |
active 顶点数的曲线怎么读,值得单独说:它是一条先冲高再递减的曲线,拐点位置反映算法的性质。广搜这类"波及范围逐层扩大"的算法,active 数会持续上升很久才回落;最短路这类"更新次数有限"的算法,通常是第一轮全图更新、之后迅速收窄到波前。曲线长期贴在高位不掉,先怀疑停机条件、再怀疑性能 —— "所有顶点每轮都发消息"是最常见的写法错误。
消息量 / 顶点数这个比值是顶点中心的第一指标:它直接决定网络流量与缓冲占用,也因此决定内存能装多少。单源最短路这类"接收方只关心最小值"的算法,启用 Combiner 后消息流量降到原来的四分之一以下。判断该不该写 Combiner 的判据很统一:消息之间做归约是否与 Compute() 的语义一致(是否可交换、可结合)。
三条判读原则:
- 看分布,不看平均值。 屏障把超步耗时定义为最慢那个 worker 的耗时,平均值漂亮但个别 worker 拖尾就足以解释超步为什么慢。先看"每个 worker 处理的顶点数与发出的消息数"的分布。
- 先分清"计算慢"还是"等待与传输慢"。 两者的解法相反:计算慢要从算法与线程数上找,等待慢要从分区与分配上找。判据是各 worker 完成的时刻是否分散 —— 分散说明进度不一,齐头并进才是真慢。
- 能从分区与算法解决的,不从加机器解决。 倾斜的根因在分区函数(哈希把高度顶点挤到了一起),加机器只会把倾斜复制到更多机器上;同理,消息量大的根因通常是没有 Combiner 或算法设计问题,加带宽只是让问题延后暴露。
最后一条与前文接得上:"顶点到机器的分配是整个模型里唯一不透明的地方",所以它也是排查时最容易漏掉的一环 —— 分区与分配在代码里是两处配置,在监控上是两条曲线。
排查时要抓的两条基线,可以让判断快很多:
- 消息总量与顶点数的比值 —— 这个比值在健康作业里是可预期的(取决于算法与是否开 Combiner)。一旦它比上一次运行翻倍,先查代码改动,参数层面解释不了这个量级的变化。
- 每超步的 active 顶点数曲线 —— 正常形态是单调递减、或先冲高再递减。出现"降下去又升回来",说明有顶点被反复唤醒,问题在消息的发送条件上,不在性能参数上。
最后一条经验:顶点中心作业的性能问题几乎都只在"分布"上显现,不在"总量"上。"总消息数偏多"这类总量指标通常是正常的(消息数天然就是边数量级),而"某个 worker 的消息数是别人的十倍"才是要修的东西。判据很直接 —— 把每个 worker 的同一指标并排画出来,看最上面那条线离中位数有多远。
相关
- RDD —— 同一个优化目标的两个表达。RDD 那篇把 Pregel 列为"能用通用抽象表达"的样本:把每一轮的顶点状态放进 RDD、用
flatMap生成消息的 RDD、再与顶点状态join回去,实现成一个 200 行的库。两篇并读能看清专用化的代价与收益:Pregel 把顶点和边留在做计算的机器上、网络只传消息,RDD 靠跨迭代一致的分区让 join 零通信 —— 后者明确说这正是 Pregel 这类专用框架的主要优化,"RDD 让用户直接表达这个目标" - GFS —— Pregel 的持久化落点:持久数据放 GFS 或 Bigtable,临时数据(缓冲的消息)放本地磁盘
- Bigtable —— 同上;图也可以存在 Bigtable 里,用行承载顶点
参考
- G. Malewicz, M. H. Austern, A. J. C. Bik, J. C. Dehnert, I. Horn, N. Leiser, G. Czajkowski. Pregel: A System for Large-Scale Graph Processing. SIGMOD 2010.
- V. Kalavri, V. Vlassov, S. Haridi. High-Level Programming Abstractions for Distributed Graph Processing. IEEE TKDE(DOI 10.1109/TKDE.2017.2762294;预印本 arXiv:1607.02646, 2016)—— 用于核对 2009–2015 年间顶点中心系统的年份、执行模型与通信机制归类
- Apache Giraph. Giraph Options. https://giraph.apache.org/options.html —— 用于核对开源实现的参数名与默认值
YJ