消息队列基础
两星期前你的接口还在每秒处理 400 个请求。今天早上它跌到每秒 30 个,日志里全是百万个没人要的推送通知造成的 ETIMEDOUT。核心服务集体崩溃,因为一个发邮件的 worker 在第三方 SMTP 接口上卡了 12 秒。这时候你需要的不是更大的数据库、也不是更好的负载均衡,而是把生产者和消费者解耦——消息队列正是实现这一点最便宜的方式,也是系统架构里最被低估的模式。
消息队列是系统两部分之间的缓冲区:产生工作的一方(生产者)和干活的一方(消费者)。生产者写一条消息,队列把它持久化保存,消费者在自己有空时拉取处理。这一个转变就消除了分布式系统里最常见的故障模式——同步耦合:一个下游服务变慢,把上游所有东西都拖垮。
为什么同步调用是根源问题
在同步架构里,每个请求都要等整条链路走完。订单服务调支付服务、支付服务调账本服务、账本再调邮件网关。如果任一环节要 3 秒,你的 p95 延迟就已经是 3 秒加开销了。当邮件网关挂了,即便你只是想发一封确认邮件,也会导致所有订单失败。

队列通过“吞掉延迟”来修复这个问题。订单服务发布一条消息,50 毫秒内就返回 200。邮件稍后发送,可能迟几秒甚至几分钟,由被允许不断重试的 worker 处理。这就是为什么支付网关普遍用 webhook 和队列,而不是把结算流水线阻塞在你的结账流程里。同理,这也是为什么健壮的后端偏好 微服务架构,而不是一条单体链路。
你必须先搞懂的核心概念
在比较工具之前,先把术语弄对,因为每家厂商对同样的东西用词略有不同:

- 生产者 / Publisher:负责把消息入队的组件。它发送完即忘,不因下游健康状态而阻塞。
- 消费者 / Subscriber:从队列读取并处理消息的组件。
- Broker(代理):负责存储和路由消息的服务端软件。RabbitMQ、Kafka、阿里云 MNS 都是 broker。
- Topic 与 Queue:Topic 把消息广播给多个消费者;Queue 把每条消息恰好交给一个消费者。Kafka 用 topic + 消费组,RabbitMQ 用 queue + 路由键。
- 至少一次 vs 恰好一次投递:几乎所有真实系统都保证“至少一次”,也就是你的 worker 可能看到同一条消息两次。消费者代码必须幂等(执行两次也安全),而不是指望 broker 去重。
- 确认应答(ACK):消费者告诉 broker 自己处理完了。如果失败或超时,broker 会重新投递。
- 死信队列(DLQ):存放反复失败消息的地方。没有 DLQ,一条“毒消息”可能被无限重试,把整条流水线卡死。
如果你跳过幂等这一课,迟早会用惨痛的方式学会:重试的支付 webhook 生成了两条扣款记录,或重试的订单任务把同一个购物车发了两次。数据处理上的这些坑,和本站中文数据分析基础里讲的去重逻辑是一脉相承的。
主要消息代理横向对比
选哪个工具,取决于你的吞吐上限、投递语义,以及你想自己运维还是租托管。

| 平台 / 工具 | 核心特点 | 价格 |
|---|---|---|
| RabbitMQ | Erlang 编写的代理,路由灵活(topic、header、direct、fanout),集群成熟,轻量且久经考验 | 开源(免费);部分云托管从 70 元/月起 |
| Apache Kafka | 分布式日志,超高吞吐(每秒百万条),可回放历史,每个分区内强有序 | 开源(免费);Confluent Cloud 免费档 5 MB/s 入站,付费约 520 元/月起 |
| Amazon SQS | 全托管,标准与 FIFO 两种模式(FIFO 有序 + 恰好一次),无服务器,不用自己跑 broker | 免费档每月 100 万次请求,之后每 100 万次约 3 元 |
| 阿里云 MNS / RocketMQ | 国内运营商友好的消息服务,与云生态集成,RocketMQ 支持事务消息 | 按量计费,起步门槛低;RocketMQ 有免费额度 |
| Apache Pulsar | 多租户、跨地域复制,计算存储分离,同时支持队列和流式语义 | 开源(免费);StreamNative 托管约 350 元/月起 |
| Redis(Streams) | 内存级速度,轻量,适合你已在用 Redis 时的简单任务队列,低量级最佳 | 开源(免费);托管 Redis 从约 70 元/月起 |
经验法则:小团队、只需要可靠的定时任务执行,用 SQS 或 Redis Streams。有多个消费者要从同一条事件流衍生不同产品(分析、搜索索引、机器学习特征)且需要回放历史时,才轮到 Kafka。国内团队如果要求数据不出境、或用阿里云生态,RocketMQ 是常见的务实选择。
怎么选:一张简单的决策树
与其照搬上一份工作用的工具,不如过一遍这四道题:

- 需要回放历史事件吗?要回放 → Kafka 或 Pulsar;不需要 → RabbitMQ 或 SQS。
- 峰值吞吐是多少?低于每秒约 5000 条 → RabbitMQ、SQS 或 Redis;更高、或大量扇出 → Kafka。
- 想自己运维基础设施吗?不想 → 托管(SQS、MNS、Confluent Cloud);选择或出于数据合规必须自建 → 自托管 RabbitMQ 或 Kafka。
- 需要严格排序 + 恰好一次吗?要 → SQS FIFO 或单分区 Kafka topic;不需要 → 至少一次就够。
团队容易过度工程化。对大多数 Web 应用处理后台邮件、文件处理、webhook 扇出,一个托管的 SQS 队列加三个 worker 容器,能撑到非常大的规模才需要 Kafka。如果你是后端新手,把 worker 跑在 WebSocket 编程入门 所讲的这类网络服务上,或借助托管基础设施把队列集成进 云端 DevOps 课程 描述的架构里,是通向可用系统最快的路。
会坑到真实团队的实现细节
只看文档不够,下面是生产环境里反复见到的坑:

- 忘了幂等键。“至少一次”意味着重复是常态而不是例外。做副作用前,先用唯一消息 ID 把已处理 ID 存进表或去重缓存。
- 把 SQS 可见超时设得太短。如果 worker 一般要 30 秒,你却设了 15 秒超时,每条慢消息都会变成重复消息。要么调高超时,要么让消费者幂等,最好两者都做。
- 没有死信配置。一次糟糕部署带来的畸形载荷会被一直重试,直到队列塞满、面向用户的系统卡死。要从第一天就接上 DLQ。
- 消费者缺少退避。崩溃的消费者堆积未 ACK 的消息,会造成重投风暴。实现指数退避和最大重试次数。
- 消费者线程里做阻塞 I/O。单线程消费者里做长数据库调用,会把整份工作串行化。并发度要匹配下游容量。
- 监控盲区。队列堆积量(等待中的消息)是最重要的指标。要在堆积超过健康基线时告警,而不是等队列空了才看。
队列系统里最常见的生产事故,不是队列挂了,而是消费者太慢、堆积增长,直到保留期到期消息被静默丢弃。用软件架构的视角看,队列只是整个系统里的一个活动部件,真正的压力点在于容量规划和监控。想打牢这块基础,可参考英文站的 数据工程基础,以及中文的 职场技能提升文章。
保留期、排序与投递语义的实战
Kafka 默认保留消息 7 天,这段时间你都得付存储费。SQS 默认最多持 14 天然后删除。RabbitMQ 一直存到被消费或设置了队列 TTL。如果你的业务需要消息保留几个月(合规、事件溯源),用 Kafka 或 Pulsar 调好保留策略,好过那种消费即删的工作队列。
排序是另一个陷阱。一个 12 分区的 topic 只保证“单个分区内”有序。如果你发布订单事件时不设分区键,相关事件可能落到不同分区、乱序到达。一定要把分区键设成你关心的实体(客户 ID、订单 ID、流 ID)。
“恰好一次”是真实但受限的。Kafka 的恰好一次语义只在“单个 Kafka 写-读-写事务”内有效,一旦碰到外部数据库或其他服务就失效。真实世界系统实际都走“至少一次 + 幂等消费者”,这也是你该为架构做的诚实规划。
值得抄走的模式
- 灾难恢复作业队列。把失败的任务(webhook、同步、导出)推到带 DLQ 和重试计划的队列,由运营人员人工清空 DLQ,而不是盯着日志看。
- 分析扇出。一个生产者事件分别喂给用户时间线服务、搜索索引和分析流水线,彼此独立。某个消费者变慢不再阻塞其余部分。
- 背压控制。因为队列吸收了尖峰,你的 worker 能保持稳定可预期的负载,不被突发压垮,顺带也平滑了数据库 CPU。
- 持久化请求/响应。即使请求-响应流程,也受益于一个请求队列加一个响应 topic。调用方轮询自己的关联 ID,中途重启不会丢请求。
你应该实际预算的成本
托管 broker 按量计费,免费档尽管慷慨但很小。SQS 每月免费 100 万次请求,算上轮询大约相当于几十万条消息,之后每 100 万次约 3 元。Confluent Cloud 免费档允许 5 MB/s 入站,够原型验证,但生产级 Kafka 的存储和出站账单会累计。谷歌 Pub/Sub 在 10 GB 免费额度后,每 TiB 投递字节收费约 280 元。国内阿里云 MNS 和 RocketMQ 入门门槛很低、有免费额度,特别适合国内团队。如果你每秒要处理上亿条消息,自托管 RabbitMQ 跑在单台虚拟机上,可能比托管方案便宜得多,代价是你得自己承担可用性和版本升级。
常见问题
什么时候该用 Kafka 而不是 RabbitMQ 或 SQS?
当你需要回放历史事件、有很多消费者读取同一条流、或吞吐超过每秒约 5000 条时,用 Kafka。标准任务执行和低量级的后台工作,用 RabbitMQ 或 SQS。
怎么阻止重复消息污染数据?
让消费者幂等:做副作用前,按消息 ID 查一张去重表,保证每条消息只处理一次。“至少一次”让重复成为常态,所以把它们当成默认,并为它们做工程设计。
死信队列是什么,什么时候该建?
DLQ 存放耗尽重试上限的消息,避免它们无限循环阻塞主队列。为你建的每一条队列都建一个 DLQ,因为一次糟糕部署带来的畸形载荷,否则会把整条流水线卡死。
实践中能拿到“恰好一次”投递吗?
只在单个 broker 事务内可以。一旦碰到外部数据库、第三方 API 或其他服务,恰好一次就不可能了。行业标准做法是“至少一次 + 幂等消费者”,你应该朝这个方向设计。
小公司跑消息队列实际要花多少钱?
非常少。SQS 免费档每月覆盖约 100 万次请求,足够处理一个典型小应用的后台负载。自托管 RabbitMQ 跑一台虚拟机,只花那台虚拟机的钱。只有当消息量冲到每月数千万条时成本才明显上升。