消息队列与批处理:吞吐与分治的底层
Kafka / RocketMQ / RabbitMQ / Pulsar 底层 + 批处理范式
本主题的底层路线(万变不离其宗):
MQ 底层 = 磁盘顺序写 + page cache + 零拷贝 + 消费位点 / 推拉模型;批处理底层 = 分治(分片) + 并行 + 流水线 + 背压。
所有吞吐与可靠,本质上都是在"顺序 IO"和"并行度"这两个旋钮上做文章:MQ 把随机写变成顺序写、把"一条条发"变成"一批批发"来获得吞吐;批处理把"一个大任务"切成"可并行的小分片"再用流水线拼起来。框架只是这套底层机制的实例化与取舍。一句话记住 先想清"顺序 IO + 并行度",四个 MQ 和三大批处理框架就只是同一套底层机制的不同脸。
0. 读前必读:四个 MQ + 批处理,其实是一套机制
你数学底子好,最怕的是被框架名词绕晕。请记住本篇的"元认知":框架是底层机制的实例化与取舍,不是新东西。我们来看这张"家族脸谱"——四种 MQ 在底层到底在争什么:
图 1:四种 MQ 的架构 / 存储取向。底层都在"写磁盘 + 取消息"上做文章,区别只是取舍点不同。
一句话记住 Kafka/RocketMQ 是"吞吐派"(顺序写+批量+拉),RabbitMQ 是"路由派"(灵活分发+单条确认),Pulsar 是"云原生派"(存算分离)——一个队列系统 80% 的差别,都来自"磁盘怎么写、消费者怎么取"。
1. MQ 通用模型:消息、通道与"不丢不重有序"
任何 MQ 都先回答三个问题:消息是什么(一段带业务语义的字节);消息去哪(通道,抽象为 topic / queue);谁消费(消费者,按模型取走消息)。在碰具体框架前,先把这套通用骨架立起来。
1.1 两种通信模型:点对点 vs 发布订阅
- 点对点(Point-to-Point / 队列):一个消息只被一个消费者取走,消费完即删除(或确认后删除)。适合"任务分发"——比如 10 台 worker 抢 1000 个任务,每个任务只做一次。
- 发布订阅(Pub/Sub):一个消息被多个订阅者同时收到,彼此独立。适合"事件广播"——比如"订单已支付"事件,库存、积分、推荐三个系统都要听。
直觉 点对点像"工单池",谁抢到谁干;发布订阅像"公众号",所有粉丝都收到。Kafka 的 consumer group 内部是点对点(组内一个分区只给一个消费者),group 之间是发布订阅(多个 group 各消费一份)。
1.2 推(push)vs 拉(pull):消费节奏谁说了算
消息从 broker 到消费者,有两种"递送动力学":
图 2:推拉模型对比。Push 延迟低但易压垮消费者;Pull(Kafka 采用)把节奏交给消费者,自带背压能力。
- Push:broker 有消息就立刻发给消费者。延迟最低,但消费者被节奏绑架——broker 不知道消费者处理得快还是慢,一窝蜂推过去就把消费者"撑死"了。需要配合消费端限速。RabbitMQ 默认偏 push。
- Pull:消费者自己按能力来取。节奏在消费者手里,天然具备"背压"(慢了就少拉点)。缺点是消费者空转时会反复来问"有吗有吗"(空轮询),一般用"长轮询"(hold 住请求直到有数据)缓解。Kafka 是典型 pull。
取舍 高吞吐场景几乎都选 pull——因为"按消费者能力取"才能稳定跑满磁盘带宽;Push 更适合"来一条立刻处理一条"的低延迟通知。
1.3 消费位点 offset:进度指针
"我消费到哪了"这个状态,在 MQ 里叫 offset(偏移量 / 位点)。它就是一个单调递增的整数指针:
- Kafka 中,每条消息在分区内有个固定 offset(0,1,2…);消费者维护"我已提交到的 offset"。
- 消费者重启后,从提交的 offset 继续,不会重复消费已提交部分,也不会漏掉——前提是 offset 提交正确。
一句话记住 offset 是"消费进度指针"。处理成功后再提交 offset = 至少一次(可能重);先提交再处理 = 至多一次(可能丢)。正确顺序决定"不丢不重"。
1.4 可靠三问:不丢 / 不重 / 有序(底层成因与对策)
消息可靠有三个经典诉求,每个都有"底层成因"和"对应机制":
| 诉求 | 底层成因(为什么会出问题) | 对策(机制) |
| 不丢(At-Least-Once) | 网络抖动 / broker 宕机 / 刷盘前断电 → 消息没落盘或没确认 | producer acks 配置、broker 多副本、flush 刷盘策略;消费端"处理完再提交 offset" |
| 不重(去重) | 网络重传、producer 重试、rebalance 重投 → 同一条被多次写入/消费 | 幂等 producer(PID+序列号去重)、消费端业务幂等(唯一键/去重表) |
| 有序 | 分区内乱序(批量、重试乱序)、跨分区无法全局有序 | 同一业务键路由到同一分区(分区内有序);单分区单线程消费 |
三个可靠等级(递进):
At-Most-Once(至多一次,可能丢,不重)→ At-Least-Once(至少一次,不丢,可能重)→ Exactly-Once(精确一次,不丢不重)。
注意:Exactly-Once 在分布式里几乎不可能"传输层面"完美达成,实际是"不丢 + 幂等去重"的工程组合,靠幂等 producer + 事务 + 位移与业务写入原子化实现。详见 §2.7。
1.5 积压与背压:生产快于消费时怎么办
当生产速率 > 消费速率,消息会在 broker 堆积(积压 / lag)。底层逻辑:
- 堆积本质:磁盘写入快(顺序写),但消费是业务逻辑,慢。缓冲区是磁盘,所以"能堆很多"——这是 MQ 解耦生产消费速率的关键价值。
- 背压(Backpressure):一种"下游扛不住就告诉上游慢点"的机制。Pull 模型天然背压(消费者自己控速);Push 模型需显式限速(如 RabbitMQ 的 prefetch count 一次只发 N 条未确认消息)。
实战 积压了别慌——先扩消费者实例(Kafka 扩 group 内消费者到分区数上限),再排查慢消费(数据库慢查询、外部调用超时)。堆积本身不会丢消息(在磁盘里),只是延迟变高。
2. Kafka:顺序写 + page cache + 零拷贝的吞吐怪兽
Kafka 是"吞吐派"鼻祖。它的全部设计都围绕一个目标:把磁盘顺序写的带宽吃满,并把读路径上的 CPU/拷贝开销压到最低。先建立核心抽象。
2.1 核心抽象:topic / partition / segment
- topic(主题):逻辑上的消息类别,如
order_paid。
- partition(分区):topic 的物理分片,分区是 Kafka 并行度的基本单位。一个 topic 有多个 partition,分布在多台 broker 上。
- segment(段):每个 partition 在磁盘上是一个只追加(append-only)的日志,被切成多个 segment 文件(如 1GB 一个),旧 segment 可删可留。
图 3:Kafka 存储与消费模型。Producer 按 key 散列写入分区,分区是只追加的日志(顺序写);消费组内每个分区仅被一个消费者消费。
2.2 存储底层:顺序写磁盘 + page cache 为什么快
这里是最该讲透的"反直觉"点:磁盘顺序写,比内存随机写还快。为什么?
- 机械盘 / SSD 的顺序吞吐都远高于随机:顺序写不需要频繁"寻道 / 移动磁头 / 随机寻址",可以一次批量刷一大片。Kafka 把"随机的业务写"转成"日志的纯顺序追加"。
- page cache(页缓存):Kafka 写消息时不直接写磁盘文件,而是写入操作系统管理的内存页缓存,由 OS 的后台线程异步刷盘。读消息也优先命中 page cache(内存)。这样"写"几乎不阻塞,"读"多数命中内存。
- 避开了"自己管内存"的坑:用 OS 的 page cache 而不是 JVM 堆——避免 GC 停顿和对象头开销,进程重启缓存还在(OS 管)。
一句话记住 Kafka 的快 = "顺序写磁盘(吃满带宽)+ 全部走 page cache(写不阻塞、读命中内存)+ 批量提交"。它不是"用内存代替磁盘",而是"把磁盘用得比内存还顺"。
2.3 发送端:批量 + 压缩 + 幂等 producer
- 批量(batching):producer 不立即发一条,而是攒一小批(
linger.ms 等几毫秒 / batch.size 攒够大小)再发。网络往返次数从"每条一次"降到"每批一次",吞吐飙升。
- 压缩:批内消息一起压缩(snappy/gzip/lz4/zstd),省网络带宽和磁盘。
- 幂等 producer(exactly-once 写入基础):producer 带一个 PID 和序列号,broker 记住"某 PID 的某序列号已写",重传的同一条会被去重。注意:幂等只保证"单分区单会话内不重",跨分区/重启仍可能重。
props.put("enable.idempotence", "true"); // 幂等 producer
props.put("acks", "all"); // 等 ISR 全确认
props.put("linger.ms", "20"); // 攒 20ms 再发,提升批量
props.put("compression.type", "lz4");
2.4 消费端:consumer group 与 rebalance
- consumer group(消费组):同一 group 内的消费者共同分担一个 topic 的所有分区,一个分区只被组内一个消费者消费(组内是点对点)。多个 group 互不干扰(组间是发布订阅)。
- rebalance(重平衡):当消费者加入/退出(扩缩容、宕机),分区归属要重新分配。期间消费暂停、offset 重新提交,有短暂停顿。频繁 rebalance 是性能杀手——用静态成员资格(static membership)可缓解。
直觉 消费并行度 ≤ 分区数。想要 10 个消费者并行,topic 至少要有 10 个分区;多于分区数的消费者会空闲。所以"分区数 = 目标并行度上限"。
2.5 副本与一致性:ISR / OSR、ack、HW / LEO、controller
Kafka 用多副本保证不丢。每个 partition 有一个 Leader 和若干 Follower,分布在不同 broker。
图 4:ISR 复制与 HW/LEO。HW(高水位)是 ISR 内最小 LEO,只有追上 HW 的消息才算"已提交、对消费者可见"。
- ISR(In-Sync Replicas,同步副本集):和 Leader 保持"差距很小"的副本集合。只有 ISR 里的副本才算数。
- OSR(Out-of-Sync Replicas):落后太多的副本被踢出 ISR。它恢复追上后才会被重新拉回。
- acks 配置(生产者确认级别):
acks=0:发了就当成功,可能丢(最快)。
acks=1:Leader 写成功就返回,Leader 宕机可能丢。
acks=-1 / all:ISR 全部写成功才返回 → 最强不丢,但延迟最高。
- LEO(Log End Offset):副本下一条要写入的位置("日志尽头")。
- HW(High Watermark,高水位):ISR 内最小的 LEO。只有 offset ≤ HW 的消息,才对消费者可见、才算"已提交"。这是防止"消费者读到还没复制完的消息,Leader 一挂消息就没了"的关键。
- controller:集群里一个特殊的 broker,负责管理分区副本的 Leader 选举、rebalance、元数据。它挂了会重新选一个,避免脑裂。
2.6 读取优化:零拷贝 sendfile
消费者读消息时,传统路径是:磁盘 → 内核缓冲区 → 用户缓冲区 → Socket 缓冲区 → 网卡,4 次拷贝 + 2 次上下文切换。Kafka 用 sendfile 系统调用(Linux 的 splice)实现零拷贝:数据直接从内核页缓存搬到网卡,不经过用户态,大幅降低 CPU 和拷贝开销。
一句话记住 零拷贝 = "数据待在内核里,从磁盘页缓存直接进网卡,不让 CPU 来回搬运"。这是 Kafka 高吞吐读的秘密武器(RabbitMQ 也用类似思路,但 Kafka/顺序读场景收益最大)。
2.7 exactly-once 思想
Kafka 的 Exactly-Once 不是"魔法传输",而是三件套:
- 幂等 producer:同一分区内去重(§2.3)。
- 事务(Transactional):producer 把"向多个分区写消息"和"提交 offset"包成一个事务,要么全成要么全不成,跨分区原子。
- 位移与业务写入原子化:消费者用事务把"处理业务"和"提交 offset"绑在一起(如写入外部 DB 和消费位移在一个事务),避免"处理了但没提交 offset(重)"或"提交了但没处理(丢)"。
清醒 "端到端 exactly-once"需要下游系统也配合(如支持事务的 sink)。Kafka 只保证自己这一段的精确一次;落到 MySQL/业务里仍需幂等兜底。
3. RocketMQ:commitlog 统一顺序写 + 业务友好的事务 / 延迟
RocketMQ(阿里开源)同样是吞吐派,但存储结构和 Kafka 不同,且更贴近"电商业务"需求(事务消息、延迟消息、死信)。
3.1 存储差异:commitlog + consumequeue + indexfile
- commitlog(核心):所有 topic 的消息全部追加写入同一个 commitlog 文件(顺序写,最大化磁盘吞吐)。这是和 Kafka 最大的区别——Kafka 是每个 partition 一个日志,RocketMQ 是所有 topic 共用一条大日志。
- consumequeue(消费队列,索引):每个 topic 的每个 queue 有一组 consumequeue 文件,里面不存消息体,只存"指向 commitlog 的偏移 + 大小 + tag"。消费者先读 consumequeue 拿到位置,再去 commitlog 取消息体。
- indexfile(索引文件):按 key / 时间建哈希索引,支持"按消息 key 查消息"。
取舍 Kafka 的"分区即文件"让单分区读很快、但 topic 多了磁盘文件碎片化;RocketMQ 的"单一 commitlog + 索引"让写入极致顺序、topic 数量不影响写性能,代价是消费要"二次查找"(先索引后取体)。
3.2 刷盘与同步:同步双写 / 异步刷盘
| 维度 | 同步刷盘(SYNC_FLUSH) | 异步刷盘(ASYNC_FLUSH) |
| 机制 | 消息写入即落盘(fsync)才返回 | 写入 page cache 即返回,后台批量刷盘 |
| 可靠 | 高(断电不丢) | 较高(极端断电可能丢缓存里几条) |
| 吞吐 | 低(等磁盘) | 高(不阻塞) |
| 主从同步 | 同步双写:主+从都写成功才返回 | 异步复制:主写成功即返回,从异步追 |
取舍 金融级要求用"同步双写 + 同步刷盘"(牺牲吞吐换不丢);普通业务用"异步"跑满吞吐。RocketMQ 让你在单条消息级别通过配置选择,比 Kafka 的全局 acks 更细。
3.3 事务消息:half 消息 + 回查
电商"下订单减库存"要保证"订单 DB 事务"和"发消息"原子。RocketMQ 的分布式事务消息流程:
- producer 先发一条 half 消息(对消费者不可见,落 commitlog)。
- producer 执行本地事务(如写订单表)。
- 成功则发 commit,half 变为可见;失败发 rollback,half 被丢弃。
- 若 broker 长时间没收到 commit/rollback(producer 宕机),会主动回查 producer "这笔事务到底成没成",再决定提交还是丢弃。
一句话记住 RocketMQ 事务消息 = "先藏一半(half,消费者看不见)→ 做本地事务 → 再决定显形或丢弃 → 失联就回查"。用"两阶段 + 回查"把本地事务和发消息绑成原子。
3.4 延迟消息与死信队列
- 延迟消息:消息发出后过 N 秒才对消费者可见(如订单 30 分钟未支付自动关闭)。RocketMQ 用预设的延迟级别(不是任意时间)实现,内部用"定时任务 + 重试队列"轮转。
- 死信队列(DLQ):消息消费重试多次仍失败(默认 16 次),被移入死信队列,不再自动重试,需人工/旁路处理。这是"失败隔离"的底层机制——坏消息不阻塞正常队列。
3.5 与 Kafka 的设计取舍
| 维度 | Kafka | RocketMQ |
| 存储 | 每分区一个日志 | 统一 commitlog + 索引 |
| 消息模型 | 都偏"拉模型 + 消费组"
| 业务特性 | 原生弱(靠外部实现事务) | 强(事务/延迟/死信开箱即用) |
| 顺序保证 | 分区内有序 | queue 内有序(可全局顺序队列) |
| 典型场景 | 日志、流处理、大吞吐 | 电商交易、业务消息、强可靠 |
4. RabbitMQ:Erlang actor 模型 + AMQP 路由
RabbitMQ 走的是"路由派"路线:不追求极致吞吐,而是把消息怎么灵活分发做到极致,底层依赖 Erlang 的并发模型。
4.1 Erlang 进程模型:轻量 actor
- actor 模型:一切皆"轻量进程(不是 OS 线程)",彼此不共享内存,只通过消息传递通信。Erlang 的进程是"廉价"的——创建成本极低、可同时存在上百万个、由 BEAM 虚拟机调度。
- 为什么适合 MQ:MQ 本质就是"海量消息在进程间流转",Erlang 的轻量 actor + 消息传递天然契合,单机就能扛大量连接与队列,且容错(supervisor 树)强。
一句话记住 RabbitMQ 的并发底层 = "Erlang 轻量 actor,不共享内存、只传消息"——和 Kafka 的"Java 线程 + 磁盘顺序写"是两条完全不同的技术路线。
4.2 AMQP 模型:exchange / binding / queue / 路由
RabbitMQ 的核心不是"topic 直接给消费者",而是经过 exchange 路由:
- Producer → Exchange(交换机):生产者只把消息发给 exchange,并带一个 routing key。
- Binding(绑定):exchange 和 queue 之间的规则(带 binding key)。
- Exchange 类型决定路由逻辑:
- direct:routing key 精确匹配 binding key(点对点式)。
- topic:routing key 按
. 分隔做模式匹配(order.*),支持发布订阅。
- fanout:广播,忽略 key,发给所有绑定队列。
- headers:按消息头属性匹配(少用)。
直觉 exchange 是"智能邮局分拣员":你只管把信(消息)交给它并写收件规则,它按规则把信投到不同邮箱(queue)。Kafka 没有这层,consumer 自己按 offset 拉。
4.3 可靠投递:publisher confirm / 镜像队列
- publisher confirm(发布确认):开启后,broker 收到消息落盘会回 ack,producer 没收到就重发 → 不丢。类似 Kafka 的 acks。
- 消息持久化:queue 和 message 都设为 durable,broker 重启不丢。
- 消费者确认(ack):消费者处理完才发 ack,broker 才删消息(
autoAck=false 防丢)。
- 镜像队列(mirrored queue):队列在多个节点存副本,主挂了从顶上(类似副本,但粒度是 queue)。新版本用 quorum queue(Raft)做更强一致。
坑 RabbitMQ 默认消息在内存队列,量大又没持久化会丢;且"推模型 + 无背压"时消费者易被冲垮,务必配 prefetch count(一次只给 N 条未确认消息,即消费端背压)。
4.4 与 Kafka / RocketMQ 的取舍
| 维度 | RabbitMQ | Kafka / RocketMQ |
| 核心取向 | 灵活路由、低延迟、单条确认 | 高吞吐、批量、日志式 |
| 并发底层 | Erlang actor | JVM 线程 + 磁盘顺序写 |
| 消息模型 | exchange 路由到 queue | topic/分区,consumer 按 offset 拉 |
| 顺序 | 单队列内有序 | 分区/queue 内有序 |
| 典型场景 | 任务分发、RPC 式调用、复杂路由 | 日志、流处理、大数据管道 |
5. Pulsar:计算存储分离 + BookKeeper 分片日志
Pulsar(雅虎开源)是"云原生派":它认为 Kafka 把"计算(broker)和存储(日志)"绑死在一台机器上,扩容麻烦。于是把两者彻底分离。
5.1 架构:Broker 无状态 + BookKeeper
- Broker(无状态):只负责接入连接、协议解析、路由,不存任何消息。挂了换一个即可,不影响数据。
- BookKeeper(存储层):专门的分布式日志存储系统,由多个 Bookie 节点组成,以"分片日志"形式持久化消息。
图 5:Pulsar 存算分离。Broker 无状态只做服务接入与路由,消息以分片日志形式落到 BookKeeper,老数据再下沉到对象存储。
5.2 分片日志 Ledger / 分层存储 / 多租户
- 分片日志(Ledger / Fragment):一个 topic 的日志被切成多个 Ledger(账本),每个 Ledger 又由多个 Bookie 上的分片组成。新消息写当前 Ledger,写满就开新的——数据天然分布在不同 Bookie,扩容只需加 Bookie。
- 分层存储(Tiered Storage):旧 Ledger 自动下沉到廉价对象存储(S3 等),新数据在 BookKeeper。兼顾"热数据低延迟 + 冷数据无限留存低成本"。
- 多租户:原生支持 tenant / namespace 层级隔离,一个集群服务多个团队互不干扰。
取舍 Pulsar 的代价是"组件多"(Broker + BookKeeper + ZooKeeper + 分层存储),运维复杂度高于 Kafka;换来的是弹性扩缩容、存储计算独立演进、多租户。
5.3 与 Kafka 架构对比
| 维度 | Kafka | Pulsar |
| 架构 | 存算一体(broker 既算又存) | 存算分离(Broker 无状态 + BookKeeper) |
| 扩容 | 扩 broker 要迁移分区数据 | 扩 Bookie 自动接管,Broker 无状态随便加 |
| 存储单元 | 分区日志(segment) | 分片 Ledger |
| 分层存储 | 需外部方案 | 原生支持 |
| 多租户 | 弱(靠不同 topic/集群) | 原生 tenant/namespace 隔离 |
一句话记住 Kafka 是"一个 broker 既服务又存盘",Pulsar 是"broker 只服务、BookKeeper 专门存"——前者简单但扩容搬数据,后者灵活但组件多。
6. 四框架横向选型对比
| 维度 | Kafka | RocketMQ | RabbitMQ | Pulsar |
| 吞吐 | 极高 | 极高 | 中 | 高 |
| 延迟 | 毫秒~秒 | 毫秒~秒 | 亚毫秒~毫秒 | 毫秒 |
| 有序 | 分区内 | queue 内 | 队列内 | 分区内 |
| 可靠 | 副本+ack | 副本+同步双写 | confirm+镜像 | BookKeeper 多副本 |
| 路由灵活 | 弱(消费组) | 弱 | 强(exchange) | 中 |
| 业务特性 | 事务弱 | 事务/延迟/死信 | 死信/延迟 | 延迟/多租户 |
| 运维 | 中 | 中 | 简单 | 复杂 |
| 典型场景 | 日志/流/大数据 | 电商/交易 | 任务/通知/路由 | 云原生/多租户 |
选型口诀 要"海量日志+流处理"选 Kafka;要"电商交易+事务延迟"选 RocketMQ;要"灵活路由+低延迟通知"选 RabbitMQ;要"弹性多租户+云原生"选 Pulsar。底层都逃不开:顺序写、副本、位点、推拉。
7. 批处理范式:分治 + 并行 + 流水线 + 背压
批处理(Batch Processing)和 MQ 是"亲戚":MQ 解决"消息一条条流",批处理解决"数据一大坨怎么算"。它们的底层共同点就是标题那套:分治(分片)+ 并行 + 流水线 + 背压。
7.1 本质:分治 + 并行 + 流水线
- 分治(Divide & Conquer):把一个大任务切成 N 个独立小任务,这是并行的前提。切出来的"片"就是 partition / split / shard。
- 并行(Parallelism):N 片分给 N 个 worker 同时算,吞吐 ≈ 单 worker × N(理想)。
- 流水线(Pipeline):不同阶段(读→处理→写)重叠执行,像工厂流水线,不让任何一段空等。
- 背压(Backpressure):下游慢了上游就减速,避免内存爆掉(和 MQ §1.5 同源)。
一句话记住 批处理 = "把大任务切成小片(分治)→ 多 worker 同时算(并行)→ 阶段重叠(流水线)→ 慢了就减速(背压)"。吞吐 = 分片数 × 并行度,瓶颈常在"跨节点搬运"(shuffle)。
7.2 数据分片(sharding)与并行流
- 分片是并行的单位:比如 1TB 文件切成 128 个 8GB 分片,每片一个任务。分片数决定最大并行度(类似 Kafka 分区数 = 消费并行度上限)。
- 并行流(Java 视角):
Stream.parallel() 或 IntStream.range().parallel() 把集合自动按 CPU 核分片并行;底层是 ForkJoinPool。适合 CPU 密集、无共享状态的本地计算。
List<Integer> r = data.parallelStream()
.map(x -> heavyCompute(x)) // 每片并行
.collect(Collectors.toList());
7.3 容错:重试 / checkpoint / 幂等
| 机制 | 底层作用 | 例子 |
| 重试(Retry) | 单片失败重跑,不丢任务 | Spring Batch retry / Flink restart |
| Checkpoint | 定期存"进度快照",失败从快照恢复而非重跑全部 | Flink barrier checkpoint / Spark RDD lineage |
| 幂等 | 重跑同片结果一致,避免重复副作用 | 唯一键写入 / 去重表 |
直觉 "重试 + checkpoint + 幂等"三件套,正是 MQ 的"不丢不重"在批处理里的翻版——分布式容错的思想完全通用。
7.4 MapReduce:分片 → map → shuffle → reduce
MapReduce(Hadoop)是批处理的"元范式",把分治流水线化。流程:
图 6:MapReduce 四段流水线:分片 → Map(并行)→ Shuffle(排序+跨节点分组,最耗资源)→ Reduce(聚合)。
- 分片(Split):输入按块切,每块一个 Map(尽量"数据本地性"——在哪台机器有数据就在哪算,省网络)。
- Map:每片并行,输出
(key, value)。可加本地 combiner 预聚合,减 Shuffle 量。
- Shuffle(洗牌):把相同 key 的数据跨节点排序、分组、搬运到对应 Reduce。这是最耗网络与磁盘的一步,也是批处理调优核心。
- Reduce:对同 key 归并聚合,写出结果。
7.5 Spring Batch:Job / Step / Reader-Processor-Writer
Spring Batch 是"企业级批处理"框架,把批处理抽象成清晰的三段式:
- Job(作业):一个完整批处理任务。
- Step(步骤):Job 由多个 Step 串成;每个 Step 是一个独立处理单元。
- Reader-Processor-Writer(读-处理-写):Step 的核心三段:
ItemReader:从 DB / 文件 / MQ 读一条数据。
ItemProcessor:处理/转换这一条。
ItemWriter:批量写出(如入库)。
- chunk(块):不是一条条提交,而是攒
chunk-size 条一次性 Writer 提交,平衡吞吐与事务。
- skip / retry / 重启:指定某些异常跳过(skip)、重试(retry);Job 有持久化"执行上下文",失败可从断点重启而非重跑全部(即 checkpoint 思想)。
@Bean
public Step step1(JobRepository repo, PlatformTransactionManager tx) {
return new StepBuilder("step1", repo)
.<In, Out>chunk(100, tx) // 每 100 条提交一次
.reader(reader()) // 读
.processor(processor()) // 处理
.writer(writer()) // 写
.faultTolerant()
.skip(SomeException.class) // 遇该异常跳过此条
.retry(OtherException.class) // 遇该异常重试
.build();
}
直觉 Spring Batch 的 Reader-Processor-Writer = "批处理版的 MQ 消费循环"(取一条→处理→确认),只是多了 Job/Step 的编排、chunk 批提交、断点重启。
7.6 Flink / Spark 批:有界流、shuffle、背压
- 有界流(Bounded):批处理的数据是"有终点"的(如昨天日志),流处理的数据是"无终点"的(如实时点击流)。Flink/Spark 统一用"算子图"表达计算。
- 算子链(Operator Chain):相邻的窄依赖算子合并成一个 task,免数据落地,提升流水线效率。
- shuffle:宽依赖(如 groupBy)需跨节点搬运数据,重新分区——就是 MR 的 Shuffle 翻版,是性能关键点。
- 背压(Backpressure):Flink 用"基于信用的流控"——下游处理慢,上游发送信用(credit)减少,自然减速,不会内存爆。Spark 用"弹性分布式数据集 RDD + 内存/磁盘溢写"缓冲。
一句话记住 Flink/Spark 批 = "把批处理画成算子 DAG,窄依赖链化提速、宽依赖 shuffle 搬运,靠背压防止下游被冲垮"。和 MQ 的推拉背压、MR 的 shuffle 是同一套思想。
7.7 流批一体:Batch = 有界流
现代框架(Flink 为主)提出流批一体:把批处理看作是"有界(有限)的数据流"。
- 统一引擎:同一套 API 和运行时,既处理无界流(实时)又处理有界流(历史批)。不再为批和流维护两套代码。
- 底层统一:批 = 有界流 = 一次性的、数据有限的流计算。Shuffle、checkpoint、状态管理全部复用。
升华 当你理解"批是有界流、MQ 是解耦的流、流处理是无界流",会发现:顺序 IO + 分片并行 + 流水线 + 背压 + checkpoint 这一套底层机制,贯穿了 MQ 和批处理两大主题。框架只是它的不同实例化。