跳转至

消息队列与异步系统

消息队列与异步系统

消息队列把生产者和消费者隔开,让任务先进入缓冲区,再由后台按能力处理。它解决的是解耦、削峰和异步,不是把失败消失。引入队列后,产品必须接受部分结果延迟,并为重复、乱序、积压和最终失败定义可见状态。

消息如何流转

flowchart LR
    producer[生产者<br>业务服务] --> broker[消息代理<br>队列或主题]
    broker --> consumer[消费者<br>Worker]
    consumer --> success{处理成功?}
    success -->|是| ack[确认消息]
    success -->|否| retry[重试或延迟重试]
    retry --> consumer
    retry -->|超过次数| dlq[死信队列]
    dlq --> repair[人工或补偿处理]

生产者只负责发布事件或任务,消费者负责处理。系统需要明确消息什么时候算「已接收」、什么时候算「已处理」,两者不能混为一个成功提示。

术语速查

术语含义产品经理要关注
生产者发布消息的服务或应用发布成功不等于下游任务完成
消费者读取并处理消息的服务或 Worker消费速度、失败重试和幂等决定最终结果
队列多条消息按规则等待消费者处理的缓冲区适合任务分发;通常一条消息由一个消费者组处理
主题 / Topic消息发布到的逻辑频道适合按事件类型组织消息,订阅者可独立消费
发布订阅一条事件广播给多个订阅者适合通知多个下游,但每个订阅者都要独立处理失败
消费者组共同消费一个队列或主题的一组实例组内扩容提高吞吐;同组通常避免重复处理同一分区消息
确认 / Ack消费者通知代理消息已处理或已接收Ack 时机决定消息丢失、重复和重试风险
至少一次消息可能重复投递,但尽量不丢消费者必须幂等,是业务系统常见的可靠性选择
至多一次消息最多投递一次,可能丢失适合丢失代价低的非关键通知或统计
恰好一次语义上只处理一次通常需要端到端约束,不能只看队列配置就宣称实现
重试处理失败后再次投递要区分临时失败和永久失败,限制次数并采用退避
死信队列超过重试上限或无法处理的消息存放处死信不能成为垃圾桶,要有告警、查看和补偿流程
幂等同一消息处理一次或多次,业务结果仍正确扣款、发券、发货、写入等操作必须有幂等键或去重记录
顺序消息按指定顺序处理相关消息顺序范围越大,吞吐和并发越受限制
背压下游处理不过来时限制上游或积累速度要定义队列容量、丢弃策略和用户可见的等待状态
消费延迟 / Lag消息产生时间与被消费时间之间的差距延迟上升是容量、故障或下游变慢的信号
事件驱动服务通过发布和订阅事件协作,而非全部同步直连降低耦合,但业务流程更难追踪和回放

队列与发布订阅

flowchart LR
    event[订单已支付]
    event --> q1[库存消费者]
    event --> q2[通知消费者]
    event --> q3[数据分析消费者]
  • 任务队列:一条任务通常由一个消费者完成,适合生成报表、发送邮件、处理文件和执行后台作业。
  • 发布订阅:一个业务事件被多个消费者独立接收,适合库存、通知、积分和分析同时响应同一事件。
  • 消费者组:组内实例分担同一份消息,组与组之间可以各自获得一份事件;扩容时要考虑分区或路由规则。

队列解决的是传递和调度,不定义业务事实。订单是否支付成功由订单或支付领域负责,消息只是把「已支付」事件传给其他系统。

同步调用与异步化

维度同步调用消息队列异步化
用户响应等待下游完成或失败先返回已提交,后续查询或通知
调用耦合上游依赖下游在线上游与消费者时间解耦
峰值处理峰值直接传到下游由队列暂存,再按消费能力处理
结果可见性当前请求可直接返回结果需要任务状态、进度、通知或对账
失败形态超时、错误码重试、积压、重复、乱序、死信
运营成本调用链较直观需要监控、回放、补偿和容量管理

适合异步化的操作有:批量导入、报表生成、文件转码、通知发送、搜索索引更新和长时间 AI 任务。支付扣款、库存锁定等需要当前请求确认结果的步骤,不能只因为响应慢就无条件改成异步。

投递语义与幂等

最常见的现实选择是至少一次投递:消费者处理成功后才确认,消费者在处理成功但确认前崩溃时,同一消息可能再次投递。因此幂等必须落在业务处理上。

sequenceDiagram
    participant Q as 消息队列
    participant C as 消费者
    participant D as 业务存储
    Q->>C: 投递消息
    C->>D: 以幂等键查询
    alt 已处理
        D-->>C: 返回已有结果
        C->>Q: Ack
    else 未处理
        C->>D: 写入结果与处理记录
        D-->>C: 提交成功
        C->>Q: Ack
    end

常见幂等做法包括唯一约束、处理记录表、幂等键、状态机和业务对账。仅在代码里加一个布尔变量不够:服务重启、并发消费和跨实例执行都可能绕过进程内状态。

重试、死信与补偿

  • 区分错误:网络超时、临时限流和服务重启可能值得重试;参数错误、权限不足和业务状态非法通常应直接进入失败处理。
  • 退避重试:逐步拉长间隔并增加随机抖动,避免大量消费者同时重试造成二次洪峰。
  • 最大次数:设置重试次数和总时间,不能让一条永久失败的消息无限占用消费资源。
  • 死信告警:死信数量、原因和消息年龄需要监控;死信页面要支持查看、修复后重放或转人工。
  • 补偿机制:下游部分成功时,通过反向操作、重试、对账或人工审核恢复业务状态。

重放消息前要确认消费者代码是否仍兼容旧格式,是否会重复发券、发通知或扣款。消息保留时间也应与审计、回放和隐私删除要求一致。

顺序、积压与容量

  • 顺序范围:按订单、用户或设备保证局部顺序,通常比全局顺序更可扩展。
  • 消费并发:并发越高吞吐越大,但可能造成同一实体乱序或数据库锁竞争。
  • 队列积压:积压量要和消息年龄一起看;数量不多但每条都很旧,同样会影响用户。
  • 背压策略:限制生产、降低非核心任务优先级、增加消费者、延长处理时间或丢弃可重建消息。
  • 容量预算:评估峰值生产速率、平均消费速率、消息大小、保留时长和重试放大倍数。

「把消费者扩容」不一定有效:分区数、数据库连接、第三方限流、单 key 热点和顺序约束都可能成为新的瓶颈。

常见黑话与真实含义

黑话需要继续追问
上 MQ 解耦了哪些服务不再同步依赖?事件格式、失败通知和流程追踪由谁负责?
消息不丢投递语义是什么?生产者、代理、消费者、数据库各段如何确认和恢复?
消费失败自动重试哪些错误可重试?重试次数、退避、死信和永久失败处理是什么?
消费者加机器就行分区、数据库、外部 API 或顺序约束是否允许并发扩容?
这个任务改异步用户如何知道任务状态?页面关闭后怎么通知?多久算失败?
消息只会消费一次是否只是队列层面一次?消费者崩溃重启后业务是否幂等?
死信后人工处理谁收到告警?如何定位原因、修复数据、重放消息并避免再次副作用?
事件最终会到最大延迟是多少?积压或消息过期时,用户和下游看到什么?

给产品经理的检查清单

  • 区分「任务已提交」「处理中」「处理成功」和「处理失败」。
  • 为每个异步任务定义状态、进度、取消、重试和通知方式。
  • 确认消息投递语义、幂等键、顺序范围和重复副作用。
  • 设定消费延迟、积压量、死信数量和消息年龄的告警阈值。
  • 为不可重试错误、永久失败、下游恢复和人工补偿设计流程。
  • 评估消息保留、数据隐私、回放兼容和容量成本。

相关页面

来源说明

本文为消息队列与异步系统通识整理,具体投递语义与实现以官方文档为准(访问日期 2026-08-30):