知识卡片
Pregel图批处理模型的顶点思维、消息传递与固定回合容错
内容
许多图算法(如PageRank)的模式是”重复沿边传播信息,直到没有更多边要跟进或某个指标收敛”,这种迭代性质普通MapReduce表达不了——它只扫一遍数据,每轮迭代都要重新读取整个输入并产出全新输出,即便相比上一轮只有极小一部分图发生了变化。Pregel模型(也叫批量同步并行/BSP模型,Apache Giraph、Spark GraphX、Flink Gelly都实现了它)专门针对这种场景优化,核心是”像顶点一样思考”:每个顶点可以沿图的边向其他顶点发送消息,框架在每次迭代里为每个顶点调用一次函数,把上一轮所有发给它的消息一并递送——如果某部分图没收到消息,那部分就不必做任何工作。与MapReduce的关键区别是顶点状态能跨迭代保留,函数只需处理新到的消息,这有点像Actor模型,但通信以固定回合方式进行:框架保证上一轮发出的所有消息,一定在本轮送达。容错通过定期把所有顶点状态存档实现,一旦某节点故障导致内存状态丢失,最简单的做法是把整个计算回滚到最近的存档点重启(若算子确定性且消息有日志记录,也可以只恢复丢失的分区)。并行执行方面,顶点不需要知道自己运行在哪台物理机器上,图的分区完全由框架决定;但由于顶点间的通信开销往往很大,若图能放进单机内存或磁盘,单机(甚至单线程)算法反而经常比分布式Pregel更快。
结构图:
flowchart TD
A[第N轮迭代开始] --> B[框架递送第N-1轮发出的全部消息]
B --> C[对每个收到消息的顶点调用处理函数]
C --> D[顶点更新自身状态,可沿边发送新消息]
D --> E{是否满足终止条件: 无更多边可跟进或指标收敛}
E -->|否| F[框架存档全部顶点状态]
F --> A
E -->|是| G[计算完成,输出结果]
参考来源
- 位置:《数据密集型应用系统设计》第十章《批处理》"图与迭代处理""Pregel处理模型""容错""并行执行"(源文件:_epub-src/ch10_split_003.html)
- 结论依据:原文说明MapReduce无法高效表达图算法的迭代传播模式,给出Pregel模型的消息传递机制、固定回合通信保证、基于存档点的容错方式,以及顶点无需知晓物理位置的并行执行方式,直接支撑本卡片的结构图与解释。
- 原始内容:一个顶点可以向另一个顶点"发送消息"……在下一轮迭代开始前,先前的迭代必须完全完成,而所有的消息必须在网络上完成复制……这种容错是通过在迭代结束时,定期存档所有顶点的状态来实现的。