跳到正文
Joeplover
学习笔记·2025-11-03·约 2 分钟阅读

WordCount 三引擎对比:Flink / MapReduce / Spark

同一个 WordCount 用 Flink、MapReduce、Spark 三种大数据引擎实现,直观对比 API 差异。

大数据对比

WordCount 三引擎对比:MapReduce vs Spark vs Flink

WordCount 是大数据领域的"Hello World"。通过同一个问题在三个不同引擎上的实现,可以直观理解它们的核心设计差异。

MapReduce 实现

public class WordCount {
    public static class TokenizerMapper 
            extends Mapper<Object, Text, Text, IntWritable> {
        public void map(Object key, Text value, Context context) throws IOException {
            StringTokenizer itr = new StringTokenizer(value.toString());
            while (itr.hasMoreTokens()) {
                word.set(itr.nextToken());
                context.write(word, one);  // (word, 1)
            }
        }
    }
    
    public static class IntSumReducer 
            extends Reducer<Text, IntWritable, Text, IntWritable> {
        public void reduce(Text key, Iterable<IntWritable> values, Context context) {
            int sum = 0;
            for (IntWritable val : values) sum += val.get();
            context.write(key, new IntWritable(sum));  // (word, count)
        }
    }
}

特点:Map 阶段写磁盘,Reduce 阶段读磁盘,中间结果强制持久化。

Spark 实现

from pyspark import SparkContext

sc = SparkContext("local", "WordCount")
text_file = sc.textFile("hdfs://input.txt")
counts = (text_file
    .flatMap(lambda line: line.split())
    .map(lambda word: (word, 1))
    .reduceByKey(lambda a, b: a + b))
counts.saveAsTextFile("hdfs://output")

特点:内存计算,DAG 优化自动合并操作(flatMap 和 map 可以流水线执行)。

public class WordCountStream {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = 
            StreamExecutionEnvironment.getExecutionEnvironment();
        
        DataStream<String> text = env.socketTextStream("localhost", 9999);
        
        text.flatMap(new Tokenizer())          // 切词
            .keyBy(value -> value.f0)          // 按 word 分组
            .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))  // 5秒窗口
            .sum(1)                            // 聚合
            .print();
        
        env.execute("Flink Streaming WordCount");
    }
}

特点:真正的流式计算,逐事件处理,毫秒级延迟。

关键对比

维度MapReduceSparkFlink
计算模型批处理批处理(底层微批次)流处理
中间结果磁盘内存内存
延迟分钟级秒级毫秒级
API 简洁度冗长简洁中等
状态管理无RDD lineage精确一次语义
速度基准1x10-100x与 Spark 相当

选型建议

  • 离线 ETL:Spark 最优。速度和开发效率的优势明显。
  • 实时计算:Flink 首选。毫秒级延迟、精确一次语义、事件时间处理。
  • 与 Hadoop 生态集成:MapReduce 最原生。但新项目不建议直接写 MR,用 Hive/Spark SQL 替代。
  • 数据分析和 ML:Spark(DataFrame/MLlib)生态最完善。

参考基准:同样的 WordCount 逻辑,Spark 比 MapReduce 快 10-100 倍(内存 vs 磁盘),Flink 在流场景下比 Spark Streaming 性能更好。