消息队列与批处理:吞吐与分治的底层

Kafka / RocketMQ / RabbitMQ / Pulsar 底层 + 批处理范式

本主题的底层路线(万变不离其宗):
MQ 底层 = 磁盘顺序写 + page cache + 零拷贝 + 消费位点 / 推拉模型;批处理底层 = 分治(分片) + 并行 + 流水线 + 背压
所有吞吐与可靠,本质上都是在"顺序 IO"和"并行度"这两个旋钮上做文章:MQ 把随机写变成顺序写、把"一条条发"变成"一批批发"来获得吞吐;批处理把"一个大任务"切成"可并行的小分片"再用流水线拼起来。框架只是这套底层机制的实例化与取舍。一句话记住 先想清"顺序 IO + 并行度",四个 MQ 和三大批处理框架就只是同一套底层机制的不同脸。
目录 0. 读前必读:四个 MQ + 批处理,其实是一套机制
1. MQ 通用模型:消息、通道与"不丢不重有序"
 1.1 两种通信模型:点对点 vs 发布订阅
 1.2 推(push)vs 拉(pull):消费节奏谁说了算
 1.3 消费位点 offset:进度指针
 1.4 可靠三问:不丢 / 不重 / 有序(底层成因与对策)
 1.5 积压与背压:生产快于消费时怎么办
2. Kafka:顺序写 + page cache + 零拷贝的吞吐怪兽
 2.1 核心抽象:topic / partition / segment
 2.2 存储底层:顺序写磁盘 + page cache 为什么快
 2.3 发送端:批量 + 压缩 + 幂等 producer
 2.4 消费端:consumer group 与 rebalance
 2.5 副本与一致性:ISR / OSR、ack、HW / LEO、controller
 2.6 读取优化:零拷贝 sendfile
 2.7 exactly-once 思想
3. RocketMQ:commitlog 统一顺序写 + 业务友好的事务 / 延迟
 3.1 存储差异:commitlog + consumequeue + indexfile
 3.2 刷盘与同步:同步双写 / 异步刷盘
 3.3 事务消息:half 消息 + 回查
 3.4 延迟消息与死信队列
 3.5 与 Kafka 的设计取舍
4. RabbitMQ:Erlang actor 模型 + AMQP 路由
 4.1 Erlang 进程模型:轻量 actor
 4.2 AMQP 模型:exchange / binding / queue / 路由
 4.3 可靠投递:publisher confirm / 镜像队列
 4.4 与 Kafka / RocketMQ 的取舍
5. Pulsar:计算存储分离 + BookKeeper 分片日志
 5.1 架构:Broker 无状态 + BookKeeper
 5.2 分片日志 Ledger / 分层存储 / 多租户
 5.3 与 Kafka 架构对比
6. 四框架横向选型对比
7. 批处理范式:分治 + 并行 + 流水线 + 背压
 7.1 本质:分治 + 并行 + 流水线
 7.2 数据分片(sharding)与并行流
 7.3 容错:重试 / checkpoint / 幂等
 7.4 MapReduce:分片 → map → shuffle → reduce
 7.5 Spring Batch:Job / Step / Reader-Processor-Writer
 7.6 Flink / Spark 批:有界流、shuffle、背压
 7.7 流批一体:Batch = 有界流

0. 读前必读:四个 MQ + 批处理,其实是一套机制

你数学底子好,最怕的是被框架名词绕晕。请记住本篇的"元认知":框架是底层机制的实例化与取舍,不是新东西。我们来看这张"家族脸谱"——四种 MQ 在底层到底在争什么:

Kafka 分区日志 顺序写+pagecache 零拷贝读 拉模型 RocketMQ commitlog 统一写 consumequeue 索引 事务/延迟消息 拉模型 RabbitMQ Erlang actor exchange 路由 AMQP 队列 推/拉混合 Pulsar 存算分离 BookKeeper 分片 分层存储 多租户 底层共性:都在"怎么把消息写进磁盘 & 怎么让消费者取"上做取舍 吞吐派(Kafka / RocketMQ):顺序写 + 批量 + 拉取 → 高吞吐、可接受秒级延迟。 路由派(RabbitMQ):灵活 exchange 路由 + 单条确认 → 低延迟、复杂路由、吞吐中。 云原生派(Pulsar):存算分离 + 分片日志 → 弹性扩缩、多租户、运维简单。
图 1:四种 MQ 的架构 / 存储取向。底层都在"写磁盘 + 取消息"上做文章,区别只是取舍点不同。
一句话记住 Kafka/RocketMQ 是"吞吐派"(顺序写+批量+拉),RabbitMQ 是"路由派"(灵活分发+单条确认),Pulsar 是"云原生派"(存算分离)——一个队列系统 80% 的差别,都来自"磁盘怎么写、消费者怎么取"。

1. MQ 通用模型:消息、通道与"不丢不重有序"

任何 MQ 都先回答三个问题:消息是什么(一段带业务语义的字节);消息去哪(通道,抽象为 topic / queue);谁消费(消费者,按模型取走消息)。在碰具体框架前,先把这套通用骨架立起来。

1.1 两种通信模型:点对点 vs 发布订阅

直觉 点对点像"工单池",谁抢到谁干;发布订阅像"公众号",所有粉丝都收到。Kafka 的 consumer group 内部是点对点(组内一个分区只给一个消费者),group 之间是发布订阅(多个 group 各消费一份)。

1.2 推(push)vs 拉(pull):消费节奏谁说了算

消息从 broker 到消费者,有两种"递送动力学":

Push(推):Broker 主动塞给消费者 Broker 有消息就推 Consumer 被动接收 优点:低延迟;缺点:消费者"被撑死"、易过载 Pull(拉):消费者自己来取 Broker 存着等取 Consumer 按需拉取 优点:消费方可控、天然背压;缺点:空轮询
图 2:推拉模型对比。Push 延迟低但易压垮消费者;Pull(Kafka 采用)把节奏交给消费者,自带背压能力。
取舍 高吞吐场景几乎都选 pull——因为"按消费者能力取"才能稳定跑满磁盘带宽;Push 更适合"来一条立刻处理一条"的低延迟通知。

1.3 消费位点 offset:进度指针

"我消费到哪了"这个状态,在 MQ 里叫 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)。底层逻辑:

实战 积压了别慌——先扩消费者实例(Kafka 扩 group 内消费者到分区数上限),再排查慢消费(数据库慢查询、外部调用超时)。堆积本身不会丢消息(在磁盘里),只是延迟变高。

2. Kafka:顺序写 + page cache + 零拷贝的吞吐怪兽

Kafka 是"吞吐派"鼻祖。它的全部设计都围绕一个目标:把磁盘顺序写的带宽吃满,并把读路径上的 CPU/拷贝开销压到最低。先建立核心抽象。

2.1 核心抽象:topic / partition / segment

Producer Topic(逻辑主题;按 key 取模散列到分区) Partition-0 segment-0.log segment-1.log 顺序追加偏移量 [offset 0..n] 顺序写磁盘 Partition-1 segment-0.log segment-1.log 顺序追加 … Partition-2 … 分区越多 并行度越高 C1 → P0 C2 → P1 消费组:每个分区只被组内一个消费者独占
图 3:Kafka 存储与消费模型。Producer 按 key 散列写入分区,分区是只追加的日志(顺序写);消费组内每个分区仅被一个消费者消费。

2.2 存储底层:顺序写磁盘 + page cache 为什么快

这里是最该讲透的"反直觉"点:磁盘顺序写,比内存随机写还快。为什么?

一句话记住 Kafka 的快 = "顺序写磁盘(吃满带宽)+ 全部走 page cache(写不阻塞、读命中内存)+ 批量提交"。它不是"用内存代替磁盘",而是"把磁盘用得比内存还顺"。

2.3 发送端:批量 + 压缩 + 幂等 producer

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

直觉 消费并行度 ≤ 分区数。想要 10 个消费者并行,topic 至少要有 10 个分区;多于分区数的消费者会空闲。所以"分区数 = 目标并行度上限"。

2.5 副本与一致性:ISR / OSR、ack、HW / LEO、controller

Kafka 用多副本保证不丢。每个 partition 有一个 Leader 和若干 Follower,分布在不同 broker。

Leader 分区 [committed] offset 0..5 offset 6,7,8(已写入) LEO=9(下条写入位置) HW=6(高水位=已提交) Follower-1(ISR) offset 0..5 offset 6,7,8 LEO=9 ≈ leader 在 ISR 集合内 Follower-2(OSR) offset 0..5 offset 6,7 LEO=8(落后) 被踢出 ISR Producer 写 Leader → 副本拉取复制 HW = ISR 内最小 LEO;只有 ≤ HW 的消息才对消费者"可见 / 可提交"。 acks=-1:等 ISR 全部确认才算提交 → 不丢;Follower 落后则被移出 ISR(OSR)。
图 4:ISR 复制与 HW/LEO。HW(高水位)是 ISR 内最小 LEO,只有追上 HW 的消息才算"已提交、对消费者可见"。

2.6 读取优化:零拷贝 sendfile

消费者读消息时,传统路径是:磁盘 → 内核缓冲区 → 用户缓冲区 → Socket 缓冲区 → 网卡,4 次拷贝 + 2 次上下文切换。Kafka 用 sendfile 系统调用(Linux 的 splice)实现零拷贝:数据直接从内核页缓存搬到网卡,不经过用户态,大幅降低 CPU 和拷贝开销。

一句话记住 零拷贝 = "数据待在内核里,从磁盘页缓存直接进网卡,不让 CPU 来回搬运"。这是 Kafka 高吞吐读的秘密武器(RabbitMQ 也用类似思路,但 Kafka/顺序读场景收益最大)。

2.7 exactly-once 思想

Kafka 的 Exactly-Once 不是"魔法传输",而是三件套:

  1. 幂等 producer:同一分区内去重(§2.3)。
  2. 事务(Transactional):producer 把"向多个分区写消息"和"提交 offset"包成一个事务,要么全成要么全不成,跨分区原子。
  3. 位移与业务写入原子化:消费者用事务把"处理业务"和"提交 offset"绑在一起(如写入外部 DB 和消费位移在一个事务),避免"处理了但没提交 offset(重)"或"提交了但没处理(丢)"。
清醒 "端到端 exactly-once"需要下游系统也配合(如支持事务的 sink)。Kafka 只保证自己这一段的精确一次;落到 MySQL/业务里仍需幂等兜底。

3. RocketMQ:commitlog 统一顺序写 + 业务友好的事务 / 延迟

RocketMQ(阿里开源)同样是吞吐派,但存储结构和 Kafka 不同,且更贴近"电商业务"需求(事务消息、延迟消息、死信)。

3.1 存储差异:commitlog + consumequeue + indexfile

取舍 Kafka 的"分区即文件"让单分区读很快、但 topic 多了磁盘文件碎片化;RocketMQ 的"单一 commitlog + 索引"让写入极致顺序、topic 数量不影响写性能,代价是消费要"二次查找"(先索引后取体)。

3.2 刷盘与同步:同步双写 / 异步刷盘

维度同步刷盘(SYNC_FLUSH)异步刷盘(ASYNC_FLUSH)
机制消息写入即落盘(fsync)才返回写入 page cache 即返回,后台批量刷盘
可靠高(断电不丢)较高(极端断电可能丢缓存里几条)
吞吐低(等磁盘)高(不阻塞)
主从同步同步双写:主+从都写成功才返回异步复制:主写成功即返回,从异步追
取舍 金融级要求用"同步双写 + 同步刷盘"(牺牲吞吐换不丢);普通业务用"异步"跑满吞吐。RocketMQ 让你在单条消息级别通过配置选择,比 Kafka 的全局 acks 更细。

3.3 事务消息:half 消息 + 回查

电商"下订单减库存"要保证"订单 DB 事务"和"发消息"原子。RocketMQ 的分布式事务消息流程:

  1. producer 先发一条 half 消息(对消费者不可见,落 commitlog)。
  2. producer 执行本地事务(如写订单表)。
  3. 成功则发 commit,half 变为可见;失败发 rollback,half 被丢弃。
  4. 若 broker 长时间没收到 commit/rollback(producer 宕机),会主动回查 producer "这笔事务到底成没成",再决定提交还是丢弃。
一句话记住 RocketMQ 事务消息 = "先藏一半(half,消费者看不见)→ 做本地事务 → 再决定显形或丢弃 → 失联就回查"。用"两阶段 + 回查"把本地事务和发消息绑成原子。

3.4 延迟消息与死信队列

3.5 与 Kafka 的设计取舍

都偏"拉模型 + 消费组"
维度KafkaRocketMQ
存储每分区一个日志统一 commitlog + 索引
消息模型
业务特性原生弱(靠外部实现事务)强(事务/延迟/死信开箱即用)
顺序保证分区内有序queue 内有序(可全局顺序队列)
典型场景日志、流处理、大吞吐电商交易、业务消息、强可靠

4. RabbitMQ:Erlang actor 模型 + AMQP 路由

RabbitMQ 走的是"路由派"路线:不追求极致吞吐,而是把消息怎么灵活分发做到极致,底层依赖 Erlang 的并发模型。

4.1 Erlang 进程模型:轻量 actor

一句话记住 RabbitMQ 的并发底层 = "Erlang 轻量 actor,不共享内存、只传消息"——和 Kafka 的"Java 线程 + 磁盘顺序写"是两条完全不同的技术路线。

4.2 AMQP 模型:exchange / binding / queue / 路由

RabbitMQ 的核心不是"topic 直接给消费者",而是经过 exchange 路由

直觉 exchange 是"智能邮局分拣员":你只管把信(消息)交给它并写收件规则,它按规则把信投到不同邮箱(queue)。Kafka 没有这层,consumer 自己按 offset 拉。

4.3 可靠投递:publisher confirm / 镜像队列

RabbitMQ 默认消息在内存队列,量大又没持久化会丢;且"推模型 + 无背压"时消费者易被冲垮,务必配 prefetch count(一次只给 N 条未确认消息,即消费端背压)。

4.4 与 Kafka / RocketMQ 的取舍

维度RabbitMQKafka / RocketMQ
核心取向灵活路由、低延迟、单条确认高吞吐、批量、日志式
并发底层Erlang actorJVM 线程 + 磁盘顺序写
消息模型exchange 路由到 queuetopic/分区,consumer 按 offset 拉
顺序单队列内有序分区/queue 内有序
典型场景任务分发、RPC 式调用、复杂路由日志、流处理、大数据管道

5. Pulsar:计算存储分离 + BookKeeper 分片日志

Pulsar(雅虎开源)是"云原生派":它认为 Kafka 把"计算(broker)和存储(日志)"绑死在一台机器上,扩容麻烦。于是把两者彻底分离

5.1 架构:Broker 无状态 + BookKeeper

Producer Consumer Broker(无状态) 只做路由/服务 不存数据 BookKeeper(存储) Bookie1 分片日志 Bookie2 分片日志 分层存储(对象存储) 旧段下沉到 S3/CFS 无限留存、成本低 计算(Broker)与存储(BookKeeper)分离:扩容 Broker 不影响数据,存储可水平扩、可分层。
图 5:Pulsar 存算分离。Broker 无状态只做服务接入与路由,消息以分片日志形式落到 BookKeeper,老数据再下沉到对象存储。

5.2 分片日志 Ledger / 分层存储 / 多租户

取舍 Pulsar 的代价是"组件多"(Broker + BookKeeper + ZooKeeper + 分层存储),运维复杂度高于 Kafka;换来的是弹性扩缩容、存储计算独立演进、多租户。

5.3 与 Kafka 架构对比

维度KafkaPulsar
架构存算一体(broker 既算又存)存算分离(Broker 无状态 + BookKeeper)
扩容扩 broker 要迁移分区数据扩 Bookie 自动接管,Broker 无状态随便加
存储单元分区日志(segment)分片 Ledger
分层存储需外部方案原生支持
多租户弱(靠不同 topic/集群)原生 tenant/namespace 隔离
一句话记住 Kafka 是"一个 broker 既服务又存盘",Pulsar 是"broker 只服务、BookKeeper 专门存"——前者简单但扩容搬数据,后者灵活但组件多。

6. 四框架横向选型对比

维度KafkaRocketMQRabbitMQPulsar
吞吐极高极高
延迟毫秒~秒毫秒~秒亚毫秒~毫秒毫秒
有序分区内queue 内队列内分区内
可靠副本+ack副本+同步双写confirm+镜像BookKeeper 多副本
路由灵活弱(消费组)强(exchange)
业务特性事务弱事务/延迟/死信死信/延迟延迟/多租户
运维简单复杂
典型场景日志/流/大数据电商/交易任务/通知/路由云原生/多租户
选型口诀 要"海量日志+流处理"选 Kafka;要"电商交易+事务延迟"选 RocketMQ;要"灵活路由+低延迟通知"选 RabbitMQ;要"弹性多租户+云原生"选 Pulsar。底层都逃不开:顺序写、副本、位点、推拉。

7. 批处理范式:分治 + 并行 + 流水线 + 背压

批处理(Batch Processing)和 MQ 是"亲戚":MQ 解决"消息一条条流",批处理解决"数据一大坨怎么算"。它们的底层共同点就是标题那套:分治(分片)+ 并行 + 流水线 + 背压

7.1 本质:分治 + 并行 + 流水线

一句话记住 批处理 = "把大任务切成小片(分治)→ 多 worker 同时算(并行)→ 阶段重叠(流水线)→ 慢了就减速(背压)"。吞吐 = 分片数 × 并行度,瓶颈常在"跨节点搬运"(shuffle)。

7.2 数据分片(sharding)与并行流

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)是批处理的"元范式",把分治流水线化。流程:

Input Split1..N (分片) Map 并行计算 (k,v) Shuffle 排序/分组 跨节点搬运 Reduce 归并聚合 (聚合) Output ① 分片:大文件切成 N 块,每块一个 Map 任务(数据本地性)。 ② Map:逐条转成 (key, value),可本地 combiner 预聚合。 ③ Shuffle:按 key 排序、跨节点分组搬运(最重网络 I/O)。 ④ Reduce:同 key 归并,输出结果。 吞吐 = 分片数 × 并行度;瓶颈常在 Shuffle 的网络与磁盘。
图 6:MapReduce 四段流水线:分片 → Map(并行)→ Shuffle(排序+跨节点分组,最耗资源)→ Reduce(聚合)。

7.5 Spring Batch:Job / Step / Reader-Processor-Writer

Spring Batch 是"企业级批处理"框架,把批处理抽象成清晰的三段式:

@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、背压

一句话记住 Flink/Spark 批 = "把批处理画成算子 DAG,窄依赖链化提速、宽依赖 shuffle 搬运,靠背压防止下游被冲垮"。和 MQ 的推拉背压、MR 的 shuffle 是同一套思想。

7.7 流批一体:Batch = 有界流

现代框架(Flink 为主)提出流批一体:把批处理看作是"有界(有限)的数据流"。

升华 当你理解"批是有界流、MQ 是解耦的流、流处理是无界流",会发现:顺序 IO + 分片并行 + 流水线 + 背压 + checkpoint 这一套底层机制,贯穿了 MQ 和批处理两大主题。框架只是它的不同实例化。
本篇为「消息队列与批处理」专题(Part 05)。底层路线:MQ = 顺序写 + page cache + 零拷贝 + 推拉/位点;批处理 = 分治 + 并行 + 流水线 + 背压。
配套 Markdown 版本见同名 .md 文件。图示建议结合 HTML 版查看(MD 版以 ASCII/文字替代并注明见 HTML 版)。