分布式计算基石:MapReduce 模型原理与 Spark 算子(Map/Reduce/Aggregate)实战指南
MapReduce is a framework for processing parallelizable problems across large datasets using a large number of computers (nodes), collectively referred to as a cluster.
使用下面的示例来简单理解下 MapReduce 的核心思想:
目标:将数据整体乘 2 后,再相加
为了理解 MapReduce 思想,就不要想先加起来,再乘以 2 的方法了~
parallelize方法将大任务分解成小任务,分发给 3 个工作节点map方法将每个节点上的数据进行..相同..的加工(每个数都乘以 2)reduce方法将加工后的数据全部加起来返回

对应 Python 代码如下
| |
| |
以上计算过程中, map 阶段..分布式..地将所有数据都乘以 2,然后 reduce 阶段将所有结果相加. 但似乎还可以继续优化, 将每个工作节点上的数据先相加起来, 再丢给 reduce 做相加计算, 会更充分地利用..分布式..的特性

具体实现时, 将 reduce 替换为 aggregate 即可
| |
本文使用了 3 种基本算子:
| 算子名称 | 基本用法 |
|---|---|
| map(func) | 映射,接收一个函数 |
| reduce(func) | 聚合,接收一个函数 |
| aggregate(zeroValue, seqOp, combOp) | zeroValue:初始值 seqOp:分区内聚合 combOp: 分区间聚合 |