高级 #distributed#mapreduce#bigdata

分布式计算(MapReduce)

MapReduce 是 Google 提出的大数据处理模型——Map 阶段并行处理数据,Shuffle 阶段按 key 分组,Reduce 阶段聚合结果。适合离线批量处理

📊 1TB 的数据——怎么在 1 小时内算完?

你有一个 1TB 的日志文件——想统计”每个页面被访问了多少次”。

单台机器:逐行读取 → 用字典计数——能算,但可能需要几小时甚至几天。

MapReduce 的思路:把数据分到 100 台机器上并行处理——然后合并结果。

🏪 **类比:全班一起数豆子”

你有一麻袋豆子(1TB 数据)——要数出”每种颜色的豆子各有多少”(统计)。

单机方式:你一个人数——数几天。

MapReduce 方式

  1. Map:100 个同学各分一堆——各人数自己这堆里每种颜色的数量
  2. Shuffle:把所有”红色”的计数集中到一起,所有”蓝色”的集中到一起……
  3. Reduce:几个同学分别把每种颜色的总数加起来

100 个同学同时干活 → 几分钟就能出结果。


🔄 MapReduce 的三个阶段

输入(大文件拆分成 M 个分片)


┌────────────────────────────────────┐
│  MAP 阶段                           │
│  输入:(文件名, 行内容)              │
│  输出:中间键值对列表                │
│  例:("page1", 1), ("page2", 1) ... │
│  每个分片由一个 Map 任务并行处理      │
└────────────────────────────────────┘


┌────────────────────────────────────┐
│  SHUFFLE 阶段(自动完成)            │
│  按 key 分组:                      │
│  "page1" → [1, 1, 1, 1]            │
│  "page2" → [1, 1]                  │
│  把相同 key 的数据发送到同一个       │
│  Reduce 任务                        │
└────────────────────────────────────┘


┌────────────────────────────────────┐
│  REDUCE 阶段                        │
│  输入:(key, value列表)             │
│  输出:(key, 聚合结果)               │
│  例:("page1", 4), ("page2", 2) ... │
└────────────────────────────────────┘


输出(结果写入分布式存储)
# MapReduce 统计页面访问量(伪代码)

def map(line):
    # 输入:日志文件的一行
    # "page1 2026-06-13 10:00:00"
    page = line.split()[0]  # 提取页面名
    emit(page, 1)            # 输出键值对:("page1", 1)

def reduce(page, counts):
    # 输入:("page1", [1, 1, 1, 1])
    total = sum(counts)
    emit(page, total)        # 输出:("page1", 4)

📐 MapReduce 的核心思想把计算推向数据,而不是把数据拉回计算。

传统做法:把 1TB 数据读到一台机器上 → 网络传输是瓶颈。 MapReduce:把计算程序分到 100 台存数据的机器上 → 每台只处理本地数据 → 网络只传中间结果。


🔧 MapReduce 之外的分布式计算

计算框架模型适用场景
MapReduce离线批处理大数据批量处理
Spark内存计算比 MapReduce 快 10-100 倍(尽量用内存)
Flink流处理实时数据处理(秒级延迟)
Pregel图计算PageRank、社交网络分析

📝 小结

概念一句话
MapReduce分而治之——Map 并行处理,Shuffle 分组,Reduce 聚合
Map输入数据 → 输出键值对(并行执行)
Shuffle按 key 分组(框架自动处理)
Reducekey + 值的列表 → 聚合结果
关键设计计算推向数据——减少网络传输
Spark内存版 MapReduce——快 10-100 倍

为什么先学这个? 批处理完成结果后,怎么通知其他系统?——消息队列(Kafka)