消息队列怎样削峰:积压、消费能力与流量整形
从生产与消费速率的差值出发,讲清消息队列如何把流量峰值转换成有完成时限的积压,以及分区、消费者并发、确认、重试、死信与背压怎样共同决定系统能否安全恢复。
活动开始后,入口在十秒内涌入几十万次报名。数据库只能稳定处理每秒八千次写入,于是架构图中间加了一条消息队列:接口收到请求后写消息,消费者再慢慢落库。图看起来合理,但它没有回答最重要的几个问题:高峰结束时会积压多少条,最后一个用户多久得到结果,队列快满时还收不收请求,消费者恢复后会不会反过来打垮数据库。
消息队列不会增加数据库的处理能力。它做的是把到达时间与处理时间分开,暂存暂时来不及处理的工作。如果峰值有限、队列装得下、业务允许稍后完成,并且高峰后的空闲能力足以清空积压,这段时间差才有价值。否则,队列只是把“接口立刻超时”改成“任务很久以后失败”。
本文回答的问题是:当入口短时流量高于下游持续处理能力时,怎样把峰值转换成有完成时限的积压,并在分区热点、重复消费、重试和消费者故障下保持容量与业务结果可控?案例继续使用活动报名系统,但方法同样适用于发短信、生成账单、异步转码、搜索索引和数据同步。
一、异步、削峰和背压不是同一件事
三个词经常一起出现,解决的问题却不同。
- 异步改变调用关系。生产者提交任务后不等待最终处理,可以先返回“已受理”。
- 削峰改变下游看到的速率。队列吸收短时差值,消费者按下游承受能力匀速处理。
- 背压把拥塞信号传回上游。系统无法继续安全接收时,要减速、拒绝或降级。
一个接口可以异步但没有削峰。消费者数量随消息量无限扩张,最终仍把峰值原样传给数据库。一个系统也可以削峰但没有有效背压:入口持续每秒写入四万条,消费者只能处理八千条,队列迟早耗尽磁盘或超过保留时间。
队列适合吸收脉冲,不适合掩盖长期容量缺口。判断能否削峰,至少要满足四个条件:峰值有明确上界或持续时间;消息能持久保存到被处理;业务允许从同步成功改成异步受理;高峰后消费能力大于新的生产速率。
二、先定义“受理成功”意味着什么
同步接口常把 HTTP 200 理解为业务已经完成。引入队列以后,接口成功可能只表示消息已被可靠接收,报名资格尚未判定。这个语义若不写进协议,用户会把“排队中”当成“报名成功”,后续任何拒绝都会变成数据事故。
一个清晰的异步接口可以返回:
HTTP/1.1 202 Accepted
Content-Type: application/json
{
"requestId": "01J...",
"status": "PENDING",
"queryUrl": "/registrations/01J..."
}
产品与技术需要共同约定:
| 契约 | 要回答的问题 |
|---|---|
| 接收保证 | 返回 202 前,消息是否已进入可恢复的持久存储 |
| 完成时限 | 99% 请求在多少秒内得到最终结果 |
| 结果获取 | 客户端轮询、回调、推送,还是后续页面查询 |
| 重复提交 | 相同用户和活动重复点击是否合并为同一请求 |
| 容量耗尽 | 拒绝新请求、只接收高优先级,还是降级到预约 |
| 最终失败 | 超过重试次数后谁处理,用户看到什么状态 |
消息写入客户端缓冲区并不等于被 broker 持久接收。Kafka 生产者可通过 acks 与副本配置选择确认强度;RabbitMQ 则用 publisher confirms 告知发布者消息是否被 broker 接管。RabbitMQ 的可靠性指南明确区分 publisher confirms 与 consumer acknowledgements:前者解决发布者到 broker,后者解决 broker 到消费者。两段都要设计,不能用一个“发送成功”概括整条链路。
如果一次请求既要写本地数据库,又要发布消息,还存在“双写”窗口。可用本地消息表或 Transactional Outbox 把业务记录与待发布事件放在同一数据库事务中,再由后台任务发送。这里不展开协议细节,可参见分布式事务:2PC、TCC、Saga 与 Outbox。
三、削峰首先是一道速率题
设时刻 t 的生产速率为 P(t),有效消费速率为 C(t),积压量为 B(t)。忽略离散误差时:
B(t) = max(0, B(0) + ∫(P(t) - C(t))dt)
当 P > C,差值进入队列;当 C > P,积压开始下降。高峰结束后的生产速率若为 P_normal,当前积压为 B_peak,理论清空时间是:
drain_time = B_peak / (C - P_normal) 前提:C > P_normal
假设活动开始后的十秒内生产速率为 40,000 条/秒,消费者受数据库保护阈值限制,只处理 8,000 条/秒:
高峰新增积压 = (40,000 - 8,000) × 10 = 320,000 条
十秒后入口回落到 2,000 条/秒,消费者仍处理 8,000 条/秒,净清空能力为 6,000 条/秒:
清空时间 = 320,000 / (8,000 - 2,000) ≈ 53.3 秒
这才是“队列能扛峰值”的可验证表述:队列至少容纳 32 万条新增消息,高峰后的约 53 秒才能完全恢复。如果产品要求全部请求 30 秒内完成,这个方案不合格;增加磁盘只会让任务保存得更久,不会缩短等待。
条数不等于容量
队列容量还要换算为字节、保留时间与副本成本。32 万条消息每条平均 2 KiB,原始数据约 625 MiB;三副本、索引、日志段和协议开销会继续放大。若消费者停机两小时,容量计算要使用两小时的入口流量,而不是一次十秒峰值。
同样,lag = 100,000 不能单独说明严重程度。每秒消费十万条时只是一秒工作量,每秒消费一千条时则是一百秒。比条数更接近用户体验的指标是最老消息年龄和预计清空时间:
estimated_drain_seconds = lag / max(consumer_rate - current_producer_rate, ε)
这也是容量规划篇所说的“队列必须有恢复预算”的延伸。需要从入口 QPS、延迟和资源预算重新推导时,可先读容量规划:从 QPS、并发与延迟推导资源预算。
四、从入口到数据库有五个控制点
一条可控的消息链路不是“生产者 → MQ → 消费者”三个方框。入口准入、发布确认、broker 存储、消费调度和下游保护分别有自己的失败方式。
1. 入口准入
网关或业务服务先检查身份、参数、幂等键与活动状态,避免把明显无效的请求写进队列。准入还应读取系统拥塞状态。当预计完成时间已经超过产品承诺,继续返回“已受理”是在制造不可兑现的承诺。
拒绝不一定意味着所有请求都返回 503。系统可以关闭非核心来源、降低单用户频率、把低优先级任务改成预约,或只允许已经获得资格令牌的用户进入。容量分配属于业务策略,broker 的磁盘水位无法替代它。
2. 可靠发布
生产者要给每个命令生成稳定的 requestId 和业务幂等键,并等待符合要求的 broker 确认。发布超时的结果通常是未知:消息可能没有到达,也可能已经持久化但确认丢失。客户端重试发布会产生重复,因此消费者仍需幂等。
生产者批量发送可提高吞吐、降低每条协议开销,但会增加等待凑批的延迟和失败重发范围。Kafka 的设计文档说明其生产者会把记录累积为批次,消费者也以批次拉取。批量大小应由延迟目标和消息尺寸约束,而不是越大越好。
3. broker 存储
broker 需要为峰值积压准备磁盘吞吐、容量、网络与副本预算。保留时间必须覆盖最长可接受停机和处置时间。监控还应包含磁盘水位、不可用分区、欠同步副本、发布确认延迟和限流时间。
当存储或内部处理跟不上发布速率,broker 会限速或拒绝。RabbitMQ 的流控说明把 flow control 定义为对发布连接的节流:下游组件无法跟上时,发布者连接会在流控状态与运行状态间切换。这个机制保护 broker,但生产者必须限制本地缓冲并把拥塞继续传回入口,否则压力只是从服务端搬到了客户端内存。
4. 消费调度
消费者从队列取出工作,控制并发、批次、预取和确认时机。它的目标不是尽快取空队列,而是在数据库、第三方接口和线程池的安全容量内尽快完成。
5. 下游保护
消费者前仍需限流、连接池隔离、超时和熔断。数据库稳定写入上限是 8,000 TPS,就不能因为积压报警而把消费者扩到每秒两万次写入。清空积压时的恢复流量同样属于高峰。
五、分区决定并行度,也制造热点
Kafka 主题被拆成多个分区。同一消费组中,一个分区同一时刻交给一个消费者实例处理,因此分区数给经典消费模型的并行度设置了上限。增加 100 个消费者但只有 32 个分区,多出的实例不会获得工作。Kafka 设计文档也说明,分区既是并行单位,也是 Kafka 提供顺序保证的边界。
顺序是按 key 购买的
若同一用户的报名、取消和再次报名必须按顺序处理,可以用 userId 作为消息 key,使同一用户进入同一分区。代价是这个 key 的工作无法由多个消费者同时处理。Kafka 不提供低成本的全局顺序;把整个活动的消息都用 activityId 作为 key,会把所有流量压进一个分区,直接丢掉大部分并行能力。
活动报名更常见的做法是按 userId 或稳定路由片分区,保证单用户命令有序;活动名额的正确性由数据库条件更新、库存状态机或专用扣减协议保证。不能为了避免超卖,把几万 QPS 全塞进一个顺序分区。
平均 lag 会隐藏热分区
64 个分区中有一个积压 20 万条,其余为空,主题平均值看起来可能尚可,那个分区上的用户却一直没有结果。监控至少要展示每分区生产速率、消费速率、lag、最老消息年龄和 key 分布。若少数大客户天然产生巨量数据,可以在业务允许的情况下把 tenantId 与桶号组合成 key;需要严格单租户顺序时,则只能为热点租户单独分流或提高单分区处理效率。
扩分区也不是无副作用的日常伸缩手段。分区增加后,默认哈希映射可能变化,同一 key 的新旧消息位于不同分区,迁移期间的顺序语义需要额外处理;分区过多还增加 broker 元数据、文件句柄、复制和恢复成本。初始分区数应结合目标吞吐、单分区基准结果、未来并行度和容错预算确定。
六、消费能力受最慢的资源约束
消费者理论吞吐可以粗略写成:
consumer_rate ≈ active_workers × batch_size / average_batch_seconds
但 active_workers 受到分区数、CPU、连接池、锁竞争和下游配额限制,average_batch_seconds 也会随数据库变慢而上升。实际容量是整条链路最小值:
C_effective = min(分区可并行能力, 消费者CPU能力, DB安全TPS, 外部依赖配额)
批量能省开销,也会放大尾延迟
批量读取、批量写库和批量确认能摊薄网络往返与事务成本。若每批等待过久,低流量时一条消息也要等到凑齐;一批中有一条坏消息,整批重试还可能重复已经成功的工作。批量大小需要同时观察吞吐、P99 处理时间、单批失败范围和内存占用。
Kafka 的 max.poll.records 限制一次 poll() 返回的记录数,max.poll.interval.ms 限制两次拉取之间允许的最长处理间隔。处理一批消息超过这个间隔,消费者会被视为失败并触发重新分配。官方的消费者配置文档建议通过减少单次记录数或增加会话间隔来匹配处理时长。更稳妥的做法通常是把拉取与受控工作池协调好,避免 poll 线程被不可控的慢任务长期阻塞。
RabbitMQ 的 prefetch 限制一个消费者尚未确认的消息数。RabbitMQ prefetch 文档说明该限制按消费者应用,0 表示不设上限。无限预取会让一个消费者占住大量消息与内存,其他消费者可能无事可做;设得过小又会因网络往返降低吞吐。预取数应接近消费者实际并发与短时缓冲需求。
扩容消费者前先看瓶颈在哪
如果 CPU 已满而数据库空闲,扩容消费者有用。如果数据库连接池、行锁或第三方配额已满,扩容只会增加排队和超时,并触发更多重试。自动扩缩容应把 lag 年龄与下游饱和度一起作为信号,并设置最大消费速率。
七、确认时机决定丢失还是重复
消费者处理一条消息大致有两个动作:产生业务副作用,以及向 broker 确认完成。两者通常无法成为同一个原子事务。
先确认,再写数据库:确认后进程崩溃,消息丢失
先写数据库,再确认:写入成功后进程崩溃,消息会再次投递
大多数业务选择第二种,即至少一次投递,再用幂等消化重复。活动报名可使用 (activity_id, user_id) 唯一约束或幂等记录表:第一次消费创建报名结果,后续重复消费读取并返回同一结果。库存扣减则要把资格写入与条件扣减放在一个本地事务中,避免重复消息重复扣库存。
BEGIN;
INSERT INTO registration(activity_id, user_id, request_id, status)
VALUES (?, ?, ?, 'SUCCESS')
ON DUPLICATE KEY UPDATE request_id = request_id;
UPDATE activity
SET remaining = remaining - 1
WHERE id = ? AND remaining > 0;
COMMIT;
示例只展示边界,真实代码还要根据 UPDATE 影响行数决定成功或售罄,并保证重复请求不会再次更新库存。幂等键的生命周期至少要覆盖消息可能被重放的时间,包括死信重新投递和灾难恢复。
“Exactly-once”只有在明确范围后才有意义。Kafka 能在其事务协议覆盖的生产与消费链路内提供特定保证,但写 MySQL、调用支付或发送短信仍是外部副作用。详细边界见超时、重试、幂等与 Exactly-once。
可见性超时也是一种租约
SQS 等队列在消费者收到消息后,会在 visibility timeout 内暂时对其他消费者隐藏消息。消费者未在期限内删除消息,消息会重新可见。AWS SQS 文档建议根据处理时长设置或动态延长该期限。
期限过短会让正常慢任务被并发执行;期限过长则在消费者崩溃后很久才重试。长任务应拆成可检查点的小步骤,或定期续租并记录进度。无论期限怎样设置,进程都可能在副作用完成后、删除消息前崩溃,所以幂等仍不可省略。
八、重试流量必须离开主消费通道
依赖抖动时,最直觉的代码是捕获异常后立刻循环重试。它会同时制造三个问题:占住分区或消费者槽位,阻塞后面的健康消息;把一次请求放大成多次下游调用;所有实例在依赖恢复瞬间一起重试,形成第二个峰值。
更可控的路径是把失败分层:
- 对短暂网络错误做少量、带抖动的进程内重试。
- 仍失败的消息进入延迟重试队列或不同延迟级别的 topic。
- 达到最大次数、格式错误或违反业务前置条件的消息进入隔离区。
- 修复原因后,以受控速率重新投递,不能一次倒回主队列。
重试消息要携带原始 requestId、首次发生时间、当前次数、最近错误分类和下一次可执行时间。退避可以采用指数增长并加入随机抖动:
delay = min(max_delay, base × 2^attempt) + random_jitter
参数错误、反序列化失败、目标记录永远不存在等永久错误不应反复撞击下游。它们属于 poison message,需要尽快隔离。AWS 关于 SQS 死信队列的建议也强调通过最大接收次数把反复失败的消息移到 DLQ,以便分析,同时避免它们继续干扰主队列。
DLQ 不是垃圾桶。每类死信要有责任人、报警阈值、查询工具、修复流程、重放审批和数据保留期限。死信长期无人处理,等于业务把失败结果藏了起来。
九、流量整形要覆盖正常、重试和恢复三股流量
消费者面对的不只有实时新消息。故障恢复后,主队列积压、延迟重试到期和人工重放可能同时出现。如果三股流量共用线程池和数据库连接,它们会相互挤占。
一种实用安排是分出独立通道与预算:
| 流量 | 目标 | 典型控制 |
|---|---|---|
| 实时消息 | 保持新请求延迟 | 保留固定消费份额 |
| 历史积压 | 在恢复时限内清空 | 限制最大 drain rate |
| 自动重试 | 等依赖恢复,避免放大 | 退避、抖动、次数上限 |
| 人工重放 | 修复旧数据 | 审批、小批、可暂停 |
假设数据库安全容量为 8,000 TPS,可以先给实时消息 5,000 TPS,历史积压 2,000 TPS,重试与重放合计 1,000 TPS。实际比例随业务调整,关键是总和不能突破下游安全线。积压清空得更慢,换来数据库不发生第二次故障。
优先级也不应只靠一个队列里的数字字段。某些 broker 的高优先级实现仍会让大消息或慢任务占住消费者,形成队头阻塞。延迟目标明显不同的任务更适合拆 topic 或 queue,使用独立消费组与资源池。例如“用户报名结果”与“报名成功后发营销短信”不该争抢同一批数据库连接。
队列满时必须作出业务决定
可以用两个水位触发动作:
- 预警水位:预计完成时间接近 SLA,暂停低优先级来源并减少非必要任务。
- 拒绝水位:磁盘、保留时间或预计完成时间越界,入口拒绝新任务或切换到预约模式。
判定条件应优先使用预计完成时间,而不是一个固定消息条数。消费能力下降一半时,同样的条数意味着两倍等待。入口拿不到精确 lag 时,可以使用控制面周期发布的拥塞等级,不必让每次请求直接查询 broker。
十、活动报名系统怎样完整落地
现在把前面的数字放进一条端到端方案。业务约束如下:峰值 40,000 请求/秒,持续十秒;之后回到 2,000 请求/秒;数据库已压测确认安全写入能力为 8,000 TPS;产品允许接口返回排队中,但要求 99% 请求在 60 秒内得到结果。
入口与消息
客户端为一次点击携带请求令牌,服务端生成 requestId,使用 (activityId, userId) 作为业务幂等键。参数、登录状态和活动时间在入口同步校验;通过后可靠发布命令,收到 broker 确认才返回 202。结果查询接口从报名结果表读取 PENDING / SUCCESS / SOLD_OUT / FAILED。
消息体只保存处理所需的稳定事实与标识,不把“当前剩余库存”这类会变化的数据复制进去:
{
"eventType": "RegistrationRequested",
"requestId": "01J...",
"activityId": 42,
"userId": 9087,
"requestedAt": "2026-10-04T12:00:00Z",
"schemaVersion": 2
}
分区与消费
主题预先建立 64 个分区,消息按 userId 哈希,使同一用户的报名与取消保持顺序。不能用唯一的 activityId 做 key,否则一个热门活动只使用一个分区。消费者批量拉取,但数据库写入并发由全局限流器限制在压测安全值以内。
消费事务先检查幂等记录,再用数据库条件更新或专门库存方案判定资格。消息排队成功只代表命令被接受,最终名额仍由权威状态决定。发送通知属于后续事件,不能让短信供应商变慢阻塞报名主链路。
验算 60 秒目标
十秒峰值结束时新增积压 32 万条,之后净清空速度为 6,000 条/秒,约 53.3 秒清空。这个数字接近 60 秒目标,尚未计算消费者重平衡、数据库慢查询和长尾消息,安全余量不足。
团队至少要做一个调整:把数据库安全能力提高并经压测确认;降低入口可接受峰值;缩短峰值持续时间;或把完成时限放宽。假设优化批量写入后稳定消费达到 10,000 条/秒,则:
高峰积压 = (40,000 - 10,000) × 10 = 300,000
高峰后清空 = 300,000 / (10,000 - 2,000) = 37.5 秒
此时才有约二十秒余量覆盖调度与尾延迟。系统还可以在预测 P99 完成时间超过 55 秒时关闭低优先级入口,在超过 60 秒时停止受理,避免承诺继续恶化。
故障与恢复
消费者停机时,入口是否继续受理由“剩余队列容量”和“预计恢复后的完成时间”共同决定。恢复后消费速率仍受数据库阈值约束。重复消息由唯一键与状态机吸收;短暂数据库错误进入延迟重试;格式损坏进入 DLQ;人工重放使用独立限额。
这套设计没有承诺所有请求成功报名。它承诺的是:已受理命令不会因单个消费者故障静默消失;同一业务请求重复执行不产生第二份结果;系统在可计算时间内给出成功、售罄或失败;超出容量时明确拒绝。
十一、监控要能回答“用户还要等多久”
仅监控 broker 存活和主题总 lag,无法判断服务是否守住完成时限。指标应沿链路分层:
| 层次 | 关键指标 |
|---|---|
| 入口 | 请求速率、受理率、拒绝率、幂等命中、202 响应延迟 |
| 生产者 | 发布速率、批次大小、确认延迟、超时、重试、本地缓冲占用 |
| broker | 每分区写入、磁盘水位、保留时间、不可用分区、欠同步副本、流控 |
| 消费者 | 每分区 lag、最老消息年龄、消费速率、处理 P99、重平衡次数 |
| 重试 | 各次数流量、下次执行延迟、DLQ 新增量、重放速率 |
| 下游 | DB TPS、连接池等待、锁等待、慢查询、错误率和资源饱和度 |
| 业务 | 排队到最终结果的 P50/P95/P99、成功率、售罄率、长期 PENDING 数 |
报警要把原因与动作关联起来。例如“最老消息超过 45 秒且数据库利用率低”可能适合扩消费者;“最老消息超过 45 秒且数据库连接池满”应先限制消费、定位数据库瓶颈;“单分区 lag 上升而其他分区正常”要查热点 key,而不是盲目扩整个消费组。
Kafka 的监控文档列出了消费者 lag、每秒消费记录数、fetch 延迟等指标。平台指标名称会变化,观测目标不变:输入多少、处理多少、差值积在哪里、最老工作等了多久、清空需要多久。
压测不能只测空队列吞吐
上线前至少覆盖以下场景:
- 按预期峰值形状灌入流量,验证最大积压和完成时限。
- 暂停一半消费者,再恢复,观察重平衡和清空过程。
- 制造一个热 key,确认单分区报警与隔离策略有效。
- 注入慢 SQL、连接池耗尽和第三方超时,确认消费限速。
- 投入永久失败消息,确认不会阻塞主分区并能进入 DLQ。
- 在业务提交后、消息确认前杀进程,验证重复消费被幂等吸收。
- 批量重放历史消息,确认实时流量仍有保留容量。
- 让 broker 达到预警水位,验证拥塞能传回入口而非撑爆生产者内存。
压测结束不能只看峰值 QPS。应保存生产与消费曲线、最大 lag、最老消息年龄、清空时间、数据库资源曲线、重复率与最终业务结果,和容量模型逐项对照。
十二、哪些请求不该进入队列
消息队列不是所有慢接口的默认答案。
| 场景 | 更合适的选择 |
|---|---|
| 调用方必须立即得到最终判定,且无法接受稍后查询 | 同步调用,配合限流与容量保障 |
| 长期输入速率高于处理能力 | 扩容或减少工作量,队列只能延后失败 |
| 工作没有持久价值,过期后再处理毫无意义 | 有界内存队列或直接丢弃,并记录指标 |
| 必须跨大量 key 保持全局顺序 | 重新审视模型;全局串行会限制吞吐 |
| 简单进程内异步,进程重启后可以丢失 | 有界线程池即可,不必引入 broker |
| 需要事件回放与多个独立订阅者 | 持久日志型消息系统更合适 |
选型还要区分任务队列与事件日志。任务通常期望被一个工作者完成;事件可被多个订阅者独立读取。Kafka、RabbitMQ、SQS 等产品能力有交集,具体选择应根据顺序、回放、路由、延迟、吞吐、运维体系和云服务约束,而不是用“高吞吐就 Kafka”一句话结束。
十三、设计检查清单
准备在系统设计中加入消息队列时,可以按下面顺序检查:
- 接口返回成功表示业务完成,还是消息已受理?
- 峰值生产速率、持续时间、正常速率分别是多少?
- 经过压测的安全消费能力是多少,最慢资源是什么?
- 最大积压、预计清空时间和最老消息 SLA 是否算过?
- broker 的磁盘、网络、副本和保留时间能否容纳最坏积压?
- 分区 key 保证了谁的顺序,会不会制造热点?
- 分区数能否支撑目标并行度,扩分区怎样处理顺序变化?
- 生产者等待了什么级别的确认,发布结果未知时怎样重试?
- 消费确认在副作用之前还是之后,重复由什么幂等约束吸收?
- 批次、prefetch、poll 间隔和下游连接池是否匹配?
- 短暂错误、永久错误和坏消息是否走不同路径?
- 实时、积压、自动重试和人工重放是否有独立预算?
- 队列接近容量或完成时限时,拥塞怎样传回入口?
- 监控能否直接回答最老任务等了多久、还要多久清空?
- 消费者故障、热分区、慢下游、DLQ 与重放是否演练过?
结语
消息队列削峰可以归结为一笔速率账:生产速率暂时超过消费速率时,差值成为积压;高峰过去后,用剩余消费能力清空。工程难点在于让这笔账始终有边界。
一个可用方案需要同时给出受理语义、最大积压、完成时限、分区方式、消费上限、确认与幂等、重试隔离、背压阈值以及故障后的恢复速率。只画一个 MQ 方框,仍然不知道系统在峰值中会怎样行动;把每个控制点和数字写清楚,队列才从“中间件”变成可验证的流量整形机制。
如果这篇文章对你有帮助