消息队列
消息队列的价值可以概括成四个词:解耦、异步、削峰、可靠投递;代价也有三个:一致性变弱、链路变长、排障变难。本文先讲清楚什么时候该用 MQ,再对比 Kafka 与 RabbitMQ 两套架构范式,最后用实测说明投递语义、幂等、顺序性与死信这些"用了才知道"的关键细节。
一、什么时候该用,什么时候不该用
| 该用 | 不该用 |
|---|---|
| 上下游处理能力不匹配,需要削峰 | 强一致要求(同一事务内的操作别拆) |
| 一次操作要触发多个下游,且允许异步(下单→发券→通知→埋点) | 需要立刻拿到结果(登录校验、库存强校验) |
| 需要重试与补偿(第三方调用失败后重投) | 数据量极小、调用简单(直接 RPC 更简单) |
| 需要广播(一个事件多个消费方) | 团队没有 MQ 运维能力(先评估代价) |
| 需要解耦发布节奏(数据变更广播给多个系统) | 只是"想让接口快一点"——先查 SQL 与缓存 |
常见误用:把同步调用强行改成异步,却不处理失败重试与幂等——结果是"接口快了,但订单丢了还查不到原因"。引入 MQ 前先想清楚:失败了怎么办?重复了怎么办?
二、两种模型
| 模型 | 语义 | 代表 |
|---|---|---|
| 队列(点对点) | 一条消息只被一个消费者处理,消费后出队 | RabbitMQ Queue、RocketMQ、ActiveMQ |
| 发布/订阅 | 一条消息被所有订阅者各处理一次 | Kafka Consumer Group、RabbitMQ fanout |
核心差异(选型时最重要的分野):
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 定位 | 消息代理(Broker 负责路由、确认、重投) | 分布式日志(消息是持久化的 append-only 日志) |
| 消费后是否删除 | 确认后删除 | 不删除,按保留期(默认 7 天)保留,可重放 |
| 消费模型 | Push(Broker 推给消费者) | Pull(消费者按 offset 拉) |
| 顺序保证 | 单队列内有序 | 分区内有序 |
| 吞吐 | 万~十万级/秒 | 十万~百万级/秒(顺序磁盘 IO + 零拷贝) |
| 延迟 | 微秒~毫秒 | 毫秒级(批量换取吞吐) |
| 典型场景 | 任务队列、RPC 异步化、复杂路由 | 日志、埋点、流处理、事件溯源 |
三、Kafka 架构要点
Topic(逻辑主题)
└─ Partition 0 [msg0][msg1][msg2][msg3]… ← 每个分区是一个有序、不可变的日志
└─ Partition 1 [msg0][msg1][msg2]…
└─ Partition 2 [msg0][msg1][msg2][msg3][msg4]…
Producer → 按 key 哈希(或轮询)决定写入哪个分区
Consumer Group → 组内每个消费者负责若干分区(一个分区只被组内一个消费者消费)
Offset → 每条消息在分区内的序号,由消费者自己维护(可回退、可重放)关键机制:
- 分区是并行度与顺序性的单位:分区越多吞吐越高,但只能保证分区内有序;
- Consumer Group 实现"队列 + 广播"两种语义:组内竞争消费(队列),组间各自独立消费(广播);
- Rebalance:消费者增减时重新分配分区,期间会短暂停止消费(避免在高峰期频繁扩缩容);
- 副本与 ISR:每个分区有多个副本,只有 ISR(In-Sync Replicas)集合内的副本才可选为 leader;
acks=all+min.insync.replicas=2才能在不丢数据和不牺牲可用性之间取得平衡; - 顺序性只在分区内成立:要按用户保序,就用
userId作为消息 key,让同一用户的消息都进同一分区(实测见 §六)。
// 生产者:至少一次语义的典型配置
props.put("acks", "all"); // 等所有 ISR 副本确认
props.put("retries", 3);
props.put("enable.idempotence", true); // 幂等生产者:避免重试产生重复消息
props.put("max.in.flight.requests.per.connection", 5);
// 消费者:手动提交 offset(处理成功后再提交)
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> r : records) {
process(r); // 业务处理(必须幂等)
}
consumer.commitSync(); // 处理完成后提交,避免"未处理就提交"导致丢消息
}四、RabbitMQ 的交换机模型
生产者不直接发到队列,而是发到 Exchange,由 Exchange 按类型和绑定规则路由:
| 类型 | 路由规则 | 场景 |
|---|---|---|
| Direct | routing key 完全匹配 | 点对点任务分发 |
| Topic | 通配符匹配(order.*、#.paid) | 按事件类型分发,最灵活 |
| Fanout | 忽略 key,广播到所有绑定队列 | 广播通知 |
| Headers | 按消息头属性匹配 | 少用(复杂且慢) |
# 生产者:发送到 topic 交换机,消息持久化
channel.exchange_declare(exchange='order.events', exchange_type='topic', durable=True)
channel.basic_publish(
exchange='order.events',
routing_key='order.paid', # 消费者可用 order.* 订阅
body=json.dumps(payload),
properties=pika.BasicProperties(delivery_mode=2) # 2 = 持久化
)
# 消费者:手动 ACK,prefetch 控制并发
channel.basic_qos(prefetch_count=10) # 一次最多推 10 条未确认消息
def callback(ch, method, properties, body):
try:
handle(body)
ch.basic_ack(delivery_tag=method.delivery_tag) # 成功后确认
except Exception:
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True) # 失败重回队列
channel.basic_consume(queue='pay.queue', on_message_callback=callback)
prefetch_count很重要:不设的话 Broker 会把队列里所有消息一次性推给消费者,导致内存爆掉与负载不均。
五、可靠性与投递语义
一条消息要经过三个环节,每个环节都可能丢:
| 环节 | 风险 | 保障手段 |
|---|---|---|
| 生产 → Broker | 网络失败、Broker 未落盘 | 发送确认(acks=all / publisher confirm)+ 失败重试 + 本地消息表 |
| Broker 存储 | 宕机丢失 | 持久化(Kafka 落盘 + 多副本;RabbitMQ 持久化队列 + 消息 delivery_mode=2) |
| Broker → 消费 | 消费失败但已确认 | 手动 ACK:业务处理成功后再提交(Kafka offset / RabbitMQ ack) |
三种投递语义:
| 语义 | 含义 | 实现 |
|---|---|---|
| 最多一次 At-most-once | 可能丢,绝不重复 | 收到就确认,失败不管(日志、埋点可接受) |
| 至少一次 At-least-once | 绝不丢,可能重复 | 处理成功再确认,失败重投 —— 主流选择,必须配套幂等 |
| 恰好一次 Exactly-once | 不丢不重 | Kafka 事务 + 幂等生产者(限 Kafka 内部流转;跨数据库仍需幂等兜底) |
六、幂等消费:至少一次的前提
实测(100 条消息,其中 30 条被重复投递,模拟 at-least-once):
无幂等:投递 130 次,扣款 1300 元,余额 8700 ← 多扣了 300 元
有幂等:投递 130 次,实际处理 100 次(拦截重复 30 条),扣款 1000 元,余额 9000// 幂等键去重:用消息里的业务唯一 ID(订单号 / 消息 ID)
const processed = new Set()
async function handle(msg) {
if (processed.has(msg.msgId)) return 'skipped' // 已处理过,直接跳过
processed.add(msg.msgId)
await deductBalance(msg.orderId, msg.amount)
return 'done'
}幂等落地的三种做法:
| 做法 | 实现 | 适用 |
|---|---|---|
| 唯一键去重表 | 用业务唯一号建唯一索引,插入成功才执行 | 最通用,推荐 |
| 数据库唯一约束 | INSERT ... ON DUPLICATE KEY UPDATE / upsert | 写入类操作 |
| 状态机 | WHERE status = 'INIT' 才更新为 PAID | 有明确业务状态流转时最优雅 |
幂等键必须是业务唯一的(订单号、流水号),不要用随机生成的消息 ID——消息重投时 ID 不变才叫重投,变了就是新消息。去重记录要有过期清理,否则表无限膨胀。
七、顺序性:只在分区内成立
实测(12 条事件按 userId 哈希分到 3 个分区,各分区并行消费):
分区 0 内顺序:op3 → op6 → op9 → op12
分区 1 内顺序:op1 → op4 → op7 → op10
分区 2 内顺序:op2 → op5 → op8 → op11
实际消费(跨分区交错):P2:u2-op2 , P0:u0-op3 , P1:u1-op1 , P0:u0-op6 , P0:u0-op9 , P2:u2-op5 ...
同一用户 u0 的 4 条事件落在分区 {0} → 分区内保序成立要点:
- 全局顺序 = 只有 1 个分区(牺牲并行度),几乎所有场景都不需要;
- 需要保序的维度(用户、订单、商品)用同一个 key,保证进同一分区;
- 消费者内不要并发处理同一分区的消息(多线程池会打乱顺序),要么单线程按序处理,要么按 key 分组并发(同 key 串行、不同 key 并行)。
八、重试、退避与死信队列
实测(20 条消息,第 7 条永久失败,最多重试 3 次并指数退避):
20 条消息:总投递 22 次(含重试)
最终失败进入 DLQ:1 条 → msg#7
其余 19 条成功;msg#7 重试 3 次后进入 DLQ,未阻塞其他消息const MAX_RETRY = 3
async function consumeWithRetry(msg) {
for (let attempt = 1; attempt <= MAX_RETRY; attempt++) {
try {
await handle(msg)
return ack(msg)
} catch (err) {
if (attempt === MAX_RETRY) return sendToDLQ(msg, err) // 重试耗尽 → 死信队列
await sleep(BASE_DELAY * 2 ** (attempt - 1)) // 指数退避
}
}
}要点:
- 必须限制重试次数:永久失败的消息如果无限重试,会堵住整个队列("毒丸消息"问题);
- 必须有退避:立即重试会把故障放大(见 微服务 的熔断与重试预算);
- 必须有死信队列(DLQ):失败消息转存到 DLQ,保留上下文供排查与人工重放;
- 区分可重试与不可重试错误:参数错误(不可重试)应直接进 DLQ,网络超时(可重试)才重试。
九、削峰:用延迟换成功率
实测(突发 1000 个请求,下游处理能力固定 50/s):
无队列:成功 108 条,失败(被拒绝)892 条,耗时 1175ms
有队列:成功 1000 条,失败 0 条,耗时 20145ms(队列峰值缓冲 1000 条)
差异:无队列丢 892 条请求;有队列零丢失,代价是耗时从 1175ms 拉长到 20145ms这就是削峰的本质:用时间换成功率。代价是延迟上升,因此不适合"必须立刻返回结果"的链路。队列长度要有监控——队列持续堆积说明消费者能力不足,需要扩容而不是继续等。
十、选型对比
| 产品 | 吞吐 | 延迟 | 顺序 | 事务消息 | 典型场景 |
|---|---|---|---|---|---|
| Kafka | 极高 | 毫秒 | 分区内 | 支持 | 日志、埋点、流处理、事件总线 |
| RabbitMQ | 中 | 低(微秒~毫秒) | 队列内 | 插件支持 | 任务队列、业务异步、复杂路由 |
| RocketMQ | 高 | 低 | 队列内 | 原生支持 | 电商交易、订单链路(阿里系) |
| Pulsar | 高 | 低 | 分区内 | 支持 | 计算存储分离、多租户 |
| Redis Stream | 中 | 极低 | 分区内 | 不支持 | 轻量场景,已有 Redis 时 |
简单决策:业务消息(需要确认、重试、路由)选 RabbitMQ/RocketMQ;数据管道(高吞吐、可重放、流处理)选 Kafka。已有 Redis 且量不大时,Redis Stream 可以省掉一个组件,但要接受它持久化与运维能力的局限(见 缓存)。
十一、检查清单
- 引入 MQ 前已想清楚:失败重试策略、幂等键、堆积告警阈值;
- 生产者开启了发送确认,失败有重投或落本地消息表;
- Broker 端开启持久化与多副本(
acks=all/ 持久化队列); - 消费者手动确认,且确认发生在业务处理成功之后;
- 消费逻辑幂等,幂等键是业务唯一号;
- 需要保序的维度用同一个 key 分区,消费者不对同 key 并发处理;
- 重试有次数上限、指数退避与死信队列;
- 队列长度与消费延迟有监控告警,堆积时有扩容预案;
- 消费者设置了并发上限(Kafka 分区数 / RabbitMQ prefetch);
- 消息有 schema 约定与版本兼容策略(加字段不改语义);
- 有消息轨迹/TraceID 便于排障(见 微服务 的可观测性);
- 有重放能力(Kafka 可按 offset 回退,RabbitMQ 可从 DLQ 重投)。
十二、常见坑速查
- 自动确认(auto-ack):消息刚推过来就确认,业务失败即丢失——必须手动确认;
- 没有幂等:at-least-once 必然重复投递,不幂等就重复扣款(实测多扣 300 元);
- 幂等键用随机 ID:重投时 ID 变了,去重失效——必须用业务唯一号;
- 无限重试:一条毒丸消息堵死整个队列——限制次数 + 退避 + DLQ;
- 不区分错误类型:参数错误也重试,白费资源——不可重试错误直接进 DLQ;
- 假设全局有序:只在分区/队列内有序,跨分区必然交错(实测见 §七);
- 分区数不够或过多:不够则并行度上不去,过多则 Rebalance 变慢、文件句柄暴涨;
- 消费者无并发上限:RabbitMQ 不设
prefetch会把消息全推过来,内存爆掉; - 消息体过大:大消息拖慢整个 Broker 并放大网络压力——只传 ID 或摘要,大字段走对象存储;
- 用 MQ 做同步 RPC:需要立刻拿到结果的链路不该异步化;
- 队列堆积不处理:只在监控里看着涨,不扩容也不降级——堆积是故障前兆;
- 消息无版本/schema 管理:上游加字段导致下游反序列化失败,全量回滚;
- 事务消息被误用:以为发了事务消息就"既不丢也不重",其实仍需消费端幂等兜底。
状态与参考
- 状态:已收录(2026-09-02,由后端领域规划清单「消息队列|Kafka 架构、消费语义、可靠投递、RabbitMQ 交换机模型、选型对比」由占位页转为正式内容)。
- 版本:削峰对比、幂等去重、分区顺序性、重试与死信队列均为 Node v22.22.1 最小实现实测(输出可直接复现);Kafka/RabbitMQ 的配置与架构为标准用法——本机未安装 Kafka/RabbitMQ,Broker 相关未做实机验证,请以运行环境为准。
- 参考:Kafka 官方文档、Kafka 设计(The Log)、RabbitMQ 教程与 AMQP 概念、RocketMQ 文档、《数据密集型应用系统设计》第 11 章 流处理。
- 阅读联动:重试、退避与熔断的配套策略见 微服务;削峰在高并发架构中的位置见 系统设计;用 Redis Stream 做轻量队列的取舍见 缓存;消息去重表的索引设计见 数据库。
下一步
- [ ] 用 Docker 起一个 Kafka(或 Redpanda)单节点,实测分区数与吞吐的关系、Consumer Group Rebalance 行为
- [ ] 用 RabbitMQ 跑通 topic 交换机 + 手动 ACK + DLQ 的完整链路,补充真实管理界面截图
- [ ] 写一篇"本地消息表 vs RocketMQ 事务消息"的分布式事务对比实验
写作规范请参阅领域概览。