知识卡片
数据流引擎用算子替代僵化的MapReduce角色
内容
把多个MapReduce作业串成工作流时,本质是把前一个作业的输出目录物化到分布式文件系统、再配置成下一个作业的输入目录,这种完全物化中间状态的方式有三个代价:必须等前驱作业全部任务完成后继任务才能启动(一个慢任务拖慢整条链路);Mapper经常是多余的,只是在读刚被Reducer写出的文件、为下一阶段的分区排序重新做准备;中间数据被复制到多台机器、写入分布式文件系统,对纯临时数据而言开销过大。Spark、Tez、Flink这类数据流引擎的核心改进是把整个工作流当成单个作业处理,用更灵活的”算子”取代僵化的Map/Reduce角色——算子间不必强制排序(除非确实需要按键分组),Mapper式的预处理逻辑可以直接并入前一个算子而不必单独起一轮,调度器由于能看到全部算子间的数据依赖关系,可以把消费数据的算子和产生数据的算子调度到同一台机器上、通过共享内存缓冲区传数据而非走网络,中间状态也倾向于留在内存或本地磁盘而非落盘到HDFS。这些优化让相同的Pig/Hive/Cascading代码只需切换执行引擎配置,就能获得明显更快的执行速度。
结构图:
flowchart LR
subgraph MapReduce工作流
J1M[作业1 Map] --> J1R[作业1 Reduce] --> HDFS1[(物化到HDFS)]
HDFS1 --> J2M[作业2 Map,常为冗余] --> J2R[作业2 Reduce]
end
subgraph 数据流引擎
OpA[算子A] -->|按需分区/排序或跳过排序| OpB[算子B]
OpB -->|共享内存缓冲区,同机调度| OpC[算子C]
OpC --> Out[(最终输出到HDFS)]
end
参考来源
- 位置:《数据密集型应用系统设计》第十章《批处理》"物化中间状态""数据流引擎"(源文件:_epub-src/ch10_split_003.html)
- 结论依据:原文列举完全物化中间状态的三个缺点(等待前驱完成/冗余Mapper/不必要的多副本写入),并说明数据流引擎把工作流当单作业处理、用算子替代Map/Reduce角色、利用局部性调度减少网络复制,直接支撑本卡片的对比结构图。
- 原始内容:将中间状态存储在分布式文件系统中意味着这些文件被复制到多个节点……它们把整个工作流作为单个作业来处理,而不是把它分解为独立的子作业……调度程序能够总览全局,知道哪里需要哪些数据,因而能够利用局部性进行优化。