Message Queues

Message Queue 怎么分发工作、怎么保证不丢不重:visibility timeout 与 ack 的时序、competing consumers、重投次数与死信队列、顺序与去重,配 SQS、RabbitMQ、Kafka 的真实默认值,外加一个邮件发送队列的算例、翻车表和面试答法。

Message queue 是一种 service-to-service 的通信方式,支持 async communication。它异步接收 producer 的消息,再投递给 consumer,让系统在高峰期也能保持稳定。

Queues 常用于大规模分布式系统的请求管理:写入 latency 在小系统里可预测,但在复杂系统里更不稳定;queue 就是 buffer + decouple 的组合。

本章讲的是 work distribution 和投递语义:一条消息怎么只交给一个 worker、worker 崩了怎么办、失败多少次算坏消息、重复和乱序怎么处理。什么时候该把请求改成异步见 异步处理模式;一条消息同时给多个系统的 fan-out 见 Publish-Subscribe;Kafka、RabbitMQ、SQS 各自怎么配见 Message Brokers。

message-queue-concepts

核心概念

  • Producer/Publisher:发送消息的一方
  • Consumer/Subscriber:读取消息的一方
  • Queue:存消息的结构(FIFO / priority / delay)
  • Broker/Queue Manager:负责路由、存储、投递、ack
  • Message:payload + metadata(headers、timestamp、priority)

message-queue-core

有约束的设计问题

一个电商的交易邮件服务:下单、发货、退款都要发邮件,平时 200 封/分钟,大促 0 点有 5 分钟 6,000 封/分钟的尖峰;邮件供应商 API 平均 400 ms,偶尔超时 10 秒;同一封订单确认邮件不能发两次,也不能因为一个地址格式错误卡住后面所有邮件。一个 worker 挂掉,它手上的邮件要有人接着发。

How it works

消息会在 queue 中保存,直到被处理并删除。一个典型流程:

  1. Producer 生成 message(payload + metadata)
  2. Producer enqueue 到 broker
  3. Broker 存储(memory / disk / replicated)
  4. Consumer dequeue 并处理
  5. Consumer 发 ack,broker 删除 message

message-queue-flow

关键在第 4 步和第 5 步之间:消息已经交给 consumer,但还没确认。所有队列都用同一个思路处理这个窗口——领取不等于删除,consumer 处理成功后显式确认,确认之前 broker 保留这条消息,超时或连接断开就交给别人。不同产品的名字不同:

产品领取后的保护期默认值超时或断开后
Amazon SQSvisibility timeout30 秒,可改;从第一次收到起最长 12 小时消息重新可见,别的 consumer 能收到
RabbitMQ手动 ack 模式下,未确认消息挂在 channel 上连接断开即重新入队;ack 超时 consumer_timeout 默认 30 分钟(4.3 起只有 quorum queue 支持)超时的 channel 被关闭,未确认消息重新入队,quorum queue 会累加投递次数
Kafka没有逐条 ack,consumer 提交 offset两次 poll() 间隔超过 max.poll.interval.ms(默认 5 分钟)即被踢出组分区交给组内别的 consumer,从上次提交的 offset 重读
Google Pub/Suback deadline10 秒,最长 600 秒消息重新投递

常见机制:

  • Visibility timeout:consumer 领取后在时间窗口内完成,超时则重新投递。处理时间不确定时,定期调用 ChangeMessageVisibility 续期,而不是一开始就设得很大——设太大,worker 真崩了,消息要很久才会被别人接手。
  • Ack / Nack:成功 ack,失败可 nack 进入 retry。
  • Retry + DLQ:多次失败进入 dead-letter queue 便于排查。

Queue types / patterns

  • Point-to-point (Work Queue):一条消息只被一个 consumer 处理
  • Pub/Sub:一条消息被多个 subscriber 消费(见 Publish-Subscribe)
  • FIFO:严格顺序,通常带去重
  • Priority Queue:按优先级处理
  • Delay/Scheduled Queue:定时投递(SQS 的消息延迟最长 15 分钟,更长的延迟要交给调度器)
  • Competing Consumers:多个 consumer 并行拉取,提升吞吐

message-queue-types

Competing consumers 的并行度上限因产品而异:SQS 和 RabbitMQ 的普通队列,加 consumer 就能加并行度;Kafka 里一个分区同一时刻只归消费组内一个 consumer,consumer 数超过分区数,多出来的只会闲着,所以 Kafka 的并行度在建 topic 时就由分区数定死了。

Delivery semantics

投递语义决定了可靠性与复杂度:

  • At-most-once:先确认再处理,可能丢,但不重复
  • At-least-once:先处理再确认,不丢但可能重复(需要 idempotency)
  • Exactly-once:最贵,通常依赖 transaction / dedupe,而且只在 broker 自己的边界内成立
  • Ordering:best-effort 或分区内顺序

实践中几乎所有业务队列都选 at-least-once,把「不重复」交给 consumer 的幂等处理。原因很简单:发邮件、扣款这类副作用发生在 broker 之外,broker 的 exactly-once 管不到。SQS 标准队列在文档里明确是至少一次投递,可能出现多份副本、偶尔乱序;SQS FIFO 队列用 MessageDeduplicationId 在 5 分钟窗口内去重,窗口之外的重复它也不管。

consumer 幂等的最小写法:

-- 处理前抢占,唯一约束保证同一条消息只有一次能插入成功
INSERT INTO sent_email (message_key, sent_at)
VALUES ('order-9921:confirmation', now())
ON CONFLICT (message_key) DO NOTHING;
-- 影响行数为 0 说明已经发过,直接 ack

message_key 要用业务键(订单号加邮件类型),不要用 broker 生成的 message ID:生产者重试时 broker 会给同一内容一个新 ID。

重投次数与死信队列

一条消息反复失败有两种原因:下游暂时坏了(等一等就好),或者消息本身坏了(地址格式错、字段缺失,重试多少次都一样)。后一种叫 poison message,不隔离就会一直占着 worker。

  • SQS:在源队列的 redrive policy 里设 maxReceiveCount,收到次数超过它就移到 DLQ。文档提醒两点:设成 1 意味着一次失败就进 DLQ,太激进;标准队列的消息过期按最初入队时间计算,所以 DLQ 的保留期要比源队列长,否则消息在 DLQ 里没待多久就过期了。保留期默认 4 天,可设 60 秒到 14 天。
  • RabbitMQ quorum queue:从 4.0 开始 delivery-limit 默认是 20,超过后丢弃或按 x-dead-letter-exchange 转入死信。经典队列没有 poison message 保护(官方特性对照表里这一项是 no),consumer 反复 nack 并 requeue 会无限循环,要在应用里自己计数。
  • Kafka:broker 没有逐条的失败计数,也没有内置 DLQ。常见做法是 consumer 自己捕获失败,写到一个重试 topic 或死信 topic,再提交 offset;否则一条坏消息会卡住整个分区。

DLQ 不是垃圾桶:要对 DLQ 深度告警,要有把修好的消息重新投回源队列的操作(SQS 提供 DLQ redrive)。

算一遍:大促 0 点的邮件队列

以下流量和耗时是本题假设。

  • 容量:尖峰 6,000 封/分钟 = 100 封/秒,平均 400 ms,Little's law 给出需要 100 × 0.4 = 40 个并发发送。每个 worker 开 10 个并发,要 4 个 worker;按供应商偶尔 10 秒超时留余量,开 8 个。
  • visibility timeout:单封最坏 10 秒超时加 3 次进程内重试的退避,最坏约 45 秒,设 60 秒。如果用 Lambda 消费 SQS,AWS 建议源队列的 visibility timeout 至少是函数超时的 6 倍,函数超时 15 秒则设 90 秒。
  • 重投与 DLQ:maxReceiveCount = 5。地址格式错的邮件第一次就会失败,最坏在 5 次领取、约 5 分钟后进 DLQ,不再占用 worker;供应商整体故障时则不该让所有消息冲进 DLQ,应该暂停消费(或由熔断器快速失败并延长可见性),而不是把正常消息当坏消息处理。
  • 积压:如果只开 4 个 worker 而供应商延迟涨到 800 ms,处理能力降到 50 封/秒,5 分钟尖峰积压 (100 − 50) × 300 = 15,000 封,尖峰过后以约 47 封/秒的净速率消化,约 5 分钟清空。按 ApproximateAgeOfOldestMessage 告警,比按队列长度告警更贴近用户感受。

Advantages

  • Scalability:高峰期 write 入 queue,消费端可横向扩展
  • Decoupling:producer/consumer 彼此不依赖,降低耦合
  • Performance:async processing 提升吞吐
  • Reliability:消息持久化 + retry + DLQ
  • Load leveling:削峰填谷,避免瞬时 overload

message-queue-advantages

Trade-offs

  • Latency vs Reliability:持久化 + retry 带来更高 latency
  • Complexity:需要监控 lag、幂等、重试策略
  • Eventual Consistency:UI/数据状态会有 delay
  • Operational cost:需要维护 broker、partition、retention

When to use

常见场景:

  • Microservices:异步服务间通信,降低同步依赖
  • Background jobs:image/video processing、email/SMS
  • Event-driven:事件分发到多个系统
  • Load spikes:请求排队,稳住 backend
  • 可靠通信:服务短暂不可用仍能投递

UI/UX 设计建议(MQ 参与时)

  • Pattern A:前台只安排任务 → 告知用户稍后完成(配合 polling / notification)
  • Pattern B:前台先做“看起来完成”的部分 → 后台异步补齐(比如 fanout timeline)

常见翻车

翻车现象修法
visibility timeout 短于处理时间同一封邮件被两个 worker 同时发timeout 大于最坏处理时间,长任务定期续期;consumer 仍要幂等
先 ack 再处理worker 重启时手上的消息静默丢失处理成功后再 ack,接受重复并去重
用 broker 的 message ID 去重生产者重试产生的重复没被挡住用业务键做幂等键
没有 DLQ 或重投上限一条坏消息无限重试,吃掉 workerSQS 设 maxReceiveCount,RabbitMQ 用 quorum queue 的 delivery limit
DLQ 保留期不长于源队列进了 DLQ 的消息没来得及排查就过期DLQ 保留期设为最长,并对 DLQ 深度告警
下游整体故障时继续消费所有正常消息都被耗尽重试次数送进 DLQ熔断下游、暂停消费,恢复后 redrive
Kafka consumer 数多于分区数扩容后吞吐不涨按目标并行度规划分区数
处理慢导致超过 max.poll.interval.msconsumer 被反复踢出组,分区来回 rebalance减小 max.poll.records(默认 500)或把慢处理移出 poll 线程

Best practices

  • Idempotency:consumer 必须可重复执行,用业务键去重
  • Retry policy:指数退避 + jitter + 最大重试次数(见 异步处理模式)
  • DLQ:隔离坏消息,避免卡住主队列
  • Backpressure:限制写入或动态扩容 consumer
  • Observability:queue length、最老消息的年龄、consumer lag、error rate、DLQ 深度
  • Security:TLS、ACL、encryption at rest

RabbitMQ vs Kafka(快速对比)

  • RabbitMQ:传统 broker,routing 灵活(exchange/queue),消费即删除,适合任务型消息
  • Kafka:分区日志(distributed log),消费不删除、按 offset 读,高吞吐、可重放,streaming/event 为主
  • 选型:任务型 + routing 用 RabbitMQ 或 SQS;大量 event、多个系统各读一遍、需要重放用 Kafka。配置细节见 Message Brokers

message-queue-rabbitmq message-queue-kafka

Backpressure

当 queue 过大,可能超过 memory 或增加磁盘 IO,导致整体性能下降。Backpressure 通过限制 queue size 或 request rate 来保持吞吐和响应时间。队列满时,client 可收到 server busy 或 HTTP 503,再配合 exponential backoff 重试。SQS 标准队列还有一个容易忽略的上限:in-flight(已领取未删除)的消息约 120,000 条,达到后短轮询会返回 OverLimit,长轮询则不再返回新消息。

面试时这样回答

  1. 先说投递语义:选 at-least-once,consumer 用业务键幂等;说明 exactly-once 管不到外部副作用。
  2. 说清领取与确认:visibility timeout 或 ack 的时序,timeout 怎么根据最坏处理时间定,长任务续期。
  3. 说清失败:重投上限、DLQ、DLQ 保留期和告警、redrive;区分坏消息和下游整体故障。
  4. 说清并行与顺序:competing consumers 怎么扩;需要顺序时按业务键分组(SQS FIFO 的 message group、Kafka 的分区键),顺序的代价是同组串行。
  5. 用 Little's law 算 worker 数和积压,并说监控看最老消息年龄。

相关章节:异步处理模式、Message Brokers、Publish-Subscribe、通知投递架构、流处理架构、Circuit breaker。

Examples / Products

一手证据