分布式系统计算:MapReduce与Spark
1. 技术分析
1.1 分布式计算概述
分布式计算将大规模计算任务分发到多个节点:
分布式计算框架
MapReduce: 批处理
Spark: 内存计算
Flink: 流式计算
核心思想:
数据分区
并行处理
结果汇总
1.2 计算框架对比
| MapReduce | 批处理 | 中 | 离线计算 |
| Spark | 批处理/流处理 | 高 | 大数据分析 |
| Flink | 流处理 | 高 | 实时计算 |
1.3 MapReduce原理
MapReduce流程
Map: 数据转换
Shuffle: 数据重分区
Reduce: 结果聚合
特点:
容错性强
可扩展
简单易用
2. 核心功能实现
2.1 MapReduce
import hashlib
class MapReduceJob:
def __init__(self, mapper, reducer):
self.mapper = mapper
self.reducer = reducer
def run(self, data, num_reducers=4):
map_output = self._map_phase(data)
shuffled = self._shuffle_phase(map_output, num_reducers)
result = self._reduce_phase(shuffled)
return result
def _map_phase(self, data):
results = []
for record in data:
results.extend(self.mapper(record))
return results
def _shuffle_phase(self, map_output, num_reducers):
partitions = [[] for _ in range(num_reducers)]
for key, value in map_output:
partition = self._get_partition(key, num_reducers)
partitions[partition].append((key, value))
return partitions
def _reduce_phase(self, partitions):
results = []
for partition in partitions:
grouped = {}
for key, value in partition:
if key not in grouped:
grouped[key] = []
grouped[key].append(value)
for key, values in grouped.items():
results.append(self.reducer(key, values))
return results
def _get_partition(self, key, num_reducers):
return int(hashlib.md5(str(key).encode()).hexdigest(), 16) % num_reducers
class WordCountJob(MapReduceJob):
def __init__(self):
super().__init__(self._mapper, self._reducer)
def _mapper(self, record):
words = record.split()
return [(word, 1) for word in words]
def _reducer(self, key, values):
return (key, sum(values))
2.2 Spark核心概念
class RDD:
def __init__(self, data, partitions=4):
self.data = data
self.partitions = partitions
def map(self, func):
return RDD([func(item) for item in self.data], self.partitions)
def flat_map(self, func):
result = []
for item in self.data:
result.extend(func(item))
return RDD(result, self.partitions)
def filter(self, func):
return RDD([item for item in self.data if func(item)], self.partitions)
def reduce(self, func):
result = self.data[0]
for item in self.data[1:]:
result = func(result, item)
return result
def group_by_key(self):
groups = {}
for key, value in self.data:
if key not in groups:
groups[key] = []
groups[key].append(value)
return RDD(list(groups.items()), self.partitions)
def reduce_by_key(self, func):
grouped = self.group_by_key()
result = []
for key, values in grouped.data:
result.append((key, self._reduce_values(values, func)))
return RDD(result, self.partitions)
def _reduce_values(self, values, func):
result = values[0]
for value in values[1:]:
result = func(result, value)
return result
class SparkContext:
def __init__(self):
pass
def parallelize(self, data, num_slices=4):
return RDD(data, num_slices)
def text_file(self, path):
with open(path, 'r') as f:
lines = f.readlines()
return RDD(lines, 4)
2.3 分布式任务调度
class TaskScheduler:
def __init__(self, workers):
self.workers = workers
self.tasks = []
def submit_task(self, task):
self.tasks.append(task)
def run(self):
results = []
while self.tasks:
task = self.tasks.pop(0)
worker = self._select_worker()
try:
result = worker.execute(task)
results.append(result)
except:
self.tasks.append(task)
return results
def _select_worker(self):
return min(self.workers, key=lambda w: w.load)
class WorkerNode:
def __init__(self, worker_id):
self.worker_id = worker_id
self.load = 0
def execute(self, task):
self.load += 1
try:
result = task.run()
return result
finally:
self.load -= 1
class Task:
def __init__(self, func, args):
self.func = func
self.args = args
def run(self):
return self.func(*self.args)
3. 性能对比
3.1 计算框架对比
| MapReduce | 高 | 中 | 高 |
| Spark | 低 | 高 | 高 |
| Flink | 低 | 很高 | 高 |
3.2 Spark vs MapReduce
| 中间结果 | 磁盘 | 内存 |
| 迭代计算 | 慢 | 快 |
| API | Java | 多语言 |
3.3 任务调度对比
| FIFO | 高 | 低 | 低 |
| Fair | 中 | 高 | 中 |
| Capacity | 中 | 中 | 中 |
4. 最佳实践
4.1 计算框架选择
def choose_computation_framework(use_case):
frameworks = {
'batch': 'Spark',
'streaming': 'Flink',
'legacy': 'MapReduce'
}
return frameworks.get(use_case, 'Spark')
class ComputationFrameworkSelector:
@staticmethod
def select(config):
frameworks = {
'mapreduce': MapReduceJob,
'spark': SparkContext
}
return frameworks[config['framework']](**config.get('params', {}))
4.2 Spark优化技巧
class SparkOptimization:
@staticmethod
def optimize(rdd):
rdd = rdd.repartition(SparkOptimization._calculate_partitions(rdd))
rdd.persist()
return rdd
@staticmethod
def _calculate_partitions(rdd):
return max(1, len(rdd.data) // 1000)
5. 总结
分布式计算是大数据处理的核心:
对比数据如下:
- Spark比MapReduce快10-100倍
- Flink是实时计算最佳选择
- Spark提供更丰富的API
- 推荐使用Spark作为默认大数据处理框架





