恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
DDIA 导读(十):批处理
首页
资讯中心
/
DDIA 导读(十):批处理
DDIA 导读(十):批处理
发布时间:2026/9/13 20:57:33
本文是《Designing Data-Intensive Applications》DDIA中文译名《数据密集型应用系统设计》第 10 章的导读。DDIA 是 Martin Kleppmann 所著的分布式系统经典本系列逐章导读把书的核心概念讲清楚。一句话主旨批处理是对一批静态数据做转换的计算模式——输入是固定的有界数据集输出是新的数据集中间可以失败重来。MapReduce 是它的经典实现但它真正的价值不在 MapReduce 这个具体框架而在把数据流组织成一连串可重跑、可独立调试的步骤这个设计思想。核心概念拆解1. 批处理的本质——纯函数式数据转换输入: 有界数据集(HDFS文件/DB表)Map: 逐条转换Shuffle: 按key分组聚合Reduce: 组内聚合输出: 新数据集批处理的关键属性有界输入处理的是已落盘的固定数据集不是无尽的流无副作用作业只读输入、写输出不改输入——所以可重跑失败重来某个 task 失败只重跑那个 task不影响其他关注吞吐非延迟分钟到小时级延迟可接受换超高吞吐批处理可重跑的前提是输出原子发布——作业完成时先写临时目录提交时原子切换这样状态与输出一致重跑不会产生部分结果。2. MapReduce 三阶段——map / shuffle / reduceReduce阶段(并行, 组内聚合)Shuffle阶段(网络传输, 按key分组)Map阶段(并行, 无通信)Mapper1: 读split1逐条转换输出(k,v)Mapper2: 读split2Mapper3: 读split3按key hash分到reducerReducer1: keyA的所有vReducer2: keyB的所有vReducer3: keyC的所有v① Map 阶段每个 Mapper 读一个输入分片split逐条记录处理输出(key, value)对。Mapper 之间无通信完全并行这是 MapReduce 可扩展的基础。② Shuffle 阶段框架自动把 Map 输出按 key 分组网络传输到对应 Reducer。Shuffle 是 MapReduce 的性能瓶颈——大量数据跨节点网络传输磁盘读写密集。③ Reduce 阶段每个 Reducer 收到同一个 key 的所有 value做聚合计数/求和/去重/拼接。Reducer 之间也无通信并行。MapReduce 的天才之处开发者只写 map 和 reduce 两个函数纯逻辑框架负责并行、容错、shuffle、调度。开发者不用管分布式细节框架自动把作业分散到几百台机器。3. Shuffle——批处理的性能瓶颈Map输出写本地磁盘按key分区网络传输到reducer节点Reducer端合并排序喂给reduce函数Shuffle 的开销Map 端写本地磁盘spill、网络跨节点传输、Reduce 端合并排序merge sort。Shuffle 数据量 性能关键。优化思路减少 shuffle 数据量——Map 端预聚合combiner、过滤早做、只 shuffle 需要的列。4. MapReduce vs Spark——为什么 Spark 更快Spark: 内存DAG内存缓存Stage1内存/可选落盘Stage2MapReduce: 每步落盘写HDFSMap→Shuffle→ReduceHDFS(磁盘)下一轮Map→Shuffle→ReduceMapReduce 的局限每一步中间结果写 HDFS磁盘。多步管道要在 HDFS 上来回读写磁盘 I/O 是瓶颈。且每步是独立作业不能优化跨步。Spark 的改进内存计算中间结果可缓存到内存RDD cache但 stage 之间 shuffle 仍写本地磁盘DAG 优化把整个管道看成有向无环图DAG优化器可以合并、重排、减少 shuffle延迟求值transformations 惰性执行遇到 action 才触发整个 DAG 优化执行选 Spark 不只是更快的 MapReduce而是内存 DAG 优化的统一批处理引擎。MapReduce 的优势是极简、极稳定、适合超大规模离线不占内存所以两者看负载形态取舍。5. 批处理中的 Join——四种策略策略机制适用Sort-Merge Join两边按 join key 排序归并大对大通用Broadcast Join小表广播到所有节点大表本地 join小对大小表能放内存Partitioned Hash Join两边按 join key 同样分区每分区本地建哈希表 join大对大分区方式一致Map 端 Join小表放分布式缓存Map 阶段直接 join无 shuffle小表很小关键点Broadcast Join 要求小表能放内存Partitioned Join 要求分区方式一致不限大小——两者都避免 shuffle 大表。Sort-Merge 是大对大时的兜底要 shuffle 排序最贵。6. 数据管道——把批处理串起来原始数据作业1: 抽取/清洗中间数据作业2: 转换/去重中间数据作业3: 导入目标系统管道设计的核心问题调度作业间依赖怎么管cronAirflow数据驱动触发中间数据每步落 HDFS 还是在线传递重跑某步失败从哪步重跑而不是从头调度器如 Airflow是批处理管道的编排层管作业依赖、重跑、失败告警。作业本身是计算层MR/Spark/脚本两层分离才能独立演进——升级调度器不影响作业逻辑。问题→方案问题——要对大批静态数据做转换规模单机扛不住。场景——输入是已落盘的有界数据集要清洗/去重/聚合/导入。方案——MapReduce 把计算拆成 map逐条转换无通信并行 shuffle按 key 分组网络传输 reduce组内聚合无通信并行开发者只写两个纯函数框架管并行/容错/调度。Shuffle 是瓶颈优化靠减少 shuffle 量combiner 预聚合、早过滤、避免大表 shuffle 的 broadcast join。Spark 用内存缓存 DAG 优化减少落盘 I/O。7. 批处理 vs 流处理——有界 vs 无界批处理处理有界数据固定数据集流处理处理无界数据持续到来的事件。批处理是等数据攒够再算流处理是来一条算一条。降低延迟边界批处理有界数据, 攒够再算延迟分钟~小时流处理无界数据, 来一条算一条延迟毫秒~秒很多系统用微批micro-batch折中——把流切成小批次处理Spark Streaming延迟在批和流之间。Mermaid批处理管道全貌编排层计算层调度调度调度作业1: 抽取/清洗作业2: 转换/去重作业3: 导入目标调度器(Airflow等)管依赖/重跑/告警原始数据中间数据中间数据目标系统下一篇第 11 章——流处理。从攒批再算转向来一条算一条。