消息队列负责保存和分发消息,Worker 是运行消费代码的进程。队列不会替业务完成任务,也通常不会自动保证外部数据库“恰好写一次”。

一次异步任务的流程

1
2
3
4
5
6
7
用户提交订单
-> API 在数据库记录订单
-> 生产者发送 order.created
-> Broker 保存并投递消息
-> Worker 拉取或接收消息
-> 执行业务并写入结果
-> 成功后 ACK

如果在业务完成前 ACK,Worker 随后崩溃会造成消息丢失;如果业务完成后、ACK 前崩溃,Broker 可能再次投递,造成重复处理。

三个 Worker 会不会都处理

取决于消费模式:

  • 竞争消费:同一消费组内,一条消息交给其中一个 Worker。
  • 广播或独立订阅:每个订阅者都收到一份消息。

扩容 Worker 可以提高并行度,但上限还受队列分片、数据库连接、下游限流和单任务耗时约束。

幂等不是简单查一次

使用稳定事件 ID,并在数据库建立唯一约束:

1
2
3
4
CREATE TABLE consumed_messages (
event_id VARCHAR(64) PRIMARY KEY,
consumed_at DATETIME NOT NULL
);

在同一事务中记录消费和修改业务:

1
2
3
4
5
6
BEGIN
-> INSERT event_id
-> 若唯一键冲突,说明已经处理,直接结束
-> 更新订单或库存
COMMIT
-> ACK

只做 SELECT 再 INSERT 仍有并发窗口,唯一索引才是最终防线。支付、发券等外部副作用还需要对方支持幂等键,或在本地维护可恢复状态机。

消息积压怎么处理

先计算生产速率与消费速率:

1
2
净积压速率 = 每秒生产数量 - 每秒消费数量
预计清空时间 = 当前 Lag / 可用净消费速率

排查顺序:

  1. 单条任务是否变慢或持续失败。
  2. 数据库、缓存和第三方接口是否成为瓶颈。
  3. Worker 是否被阻塞、崩溃或频繁重启。
  4. 队列分片是否足以支持更多消费者。
  5. 扩容是否会压垮下游。

Webhook 与 MQ 的区别

Webhook 是系统之间通过 HTTP 主动通知;MQ 是通过 Broker 缓冲和分发消息。Webhook 也需要重试、签名、幂等和失败补偿,不能因为使用 HTTP 就省略可靠性设计。

延伸阅读

站内搜索

没有找到内容!