消息队列架构踩坑记录

最早那版系统是同步调用:用户一点下单,接口里连着发短信、推邮件、扣库存、写日志,任何一个下游抖一下,整条链路都跟着卡。

接口 P99 latency 被拖到了不可接受的程度,我们才被迫把"发通知"和"核心下单"拆开——消息队列不是架构炫技,是响应时间逼出来的。

为什么需要消息队列

最初的项目是同步调用架构:

# 用户下单
def create_order(user_id, items):
    # 创建订单
    order = Order.create(user_id, items)

    # 发送短信通知
    send_sms(order.user.phone, "订单创建成功")

    # 推送物流信息
    push_logistics(order)

    # 更新库存
    update_inventory(items)

    # 更新积分
    update_points(user_id, order.total)

    return order

问题很明显:

  • 如果短信接口慢,整个下单就慢
  • 如果物流服务挂了,下单就失败
  • 这些非核心功能阻塞了核心流程

引入消息队列后:

sequenceDiagram participant User as 用户 participant Order as 订单服务 participant MQ as 消息队列 participant Sms as 短信服务 participant Logistics as 物流服务 User->>Order: 创建订单 Order->>Order: 核心业务逻辑 Order->>MQ: 发送消息 MQ-->>User: 立即返回 Note over MQ,Logistics: 异步处理 MQ->>Sms: 推送短信任务 MQ->>Logistics: 推送物流任务

好处是立竿见影的:

  • 解耦:服务之间不直接依赖
  • 异步:主流程快速返回
  • 削峰:高峰流量缓存处理

消息模型选择

消息队列有两种基本模型:

点对点模型

一个消息只能被一个消费者消费,适合任务分发。

graph TB A[生产者1] --> C[消息队列] B[生产者2] --> C C --> D[消费者1] C --> E[消费者2] C --> F[消费者3] Note over A,F: 每个消息只被一个消费者消费 style C fill:#87CEEB

典型场景:邮件发送、图片处理

发布订阅模型

一个消息可以被多个消费者消费,适合事件通知。

graph TB A[生产者] --> B[交换机/Topic] B --> C[队列1-短信] B --> D[队列2-邮件] B --> E[队列3-日志] C --> F[消费者1] D --> G[消费者2] E --> H[消费者3] Note over A,H: 每个消息被多个消费者消费 style B fill:#FFD700

典型场景:订单事件、系统通知

可靠性保证机制

消息持久化

消息队列最重要的特性是可靠性。如果消息丢失,业务就出问题了。

持久化机制

  • 消息接收到后立即写入磁盘
  • 使用同步写(fsync)确保数据真正落盘
  • 通过副本机制提供冗余

确认机制

生产者确认:生产者可以请求确认,确保消息成功写入队列。

消费者确认:消费者处理完消息后需要确认,消息队列才会删除消息。

sequenceDiagram participant Producer as 生产者 participant MQ as 消息队列 participant Consumer as 消费者 Producer->>MQ: 发送消息 MQ->>MQ: 写入磁盘 MQ-->>Producer: 确认接收 MQ->>Consumer: 推送消息 Consumer->>Consumer: 处理消息 Consumer->>MQ: 确认消息 MQ->>MQ: 删除消息

重试与死信队列

消息处理失败时,消息队列会自动重试。如果重试多次仍然失败,消息会被转移到死信队列。

graph TB A[消息处理] --> B{处理成功?} B -->|是| C[确认并删除] B -->|否| D{达到最大<br/>重试次数?} D -->|否| E[延迟重试] D -->|是| F[死信队列] E --> A style C fill:#90EE90 style F fill:#FFB6C1

选型:Kafka vs RabbitMQ

Kafka 的优势

高吞吐量:单机每秒可处理百万级消息

分布式架构:天然支持水平扩展

消息持久化:消息持久化到磁盘,支持消息回溯

流处理:与流处理框架深度集成

适合场景:日志收集、流数据处理、大数据传输

RabbitMQ 的优势

消息路由:支持复杂的路由规则和消息模式

协议标准:基于 AMQP 协议,标准化程度高

管理界面:管理界面友好,运维简单

事务支持:支持事务消息,保证强一致性

适合场景:复杂路由、企业集成、分布式事务

我的选型

项目特点是:

  • 流量中等,不需要 Kafka 的高吞吐量
  • 需要复杂路由,不同消息走不同处理逻辑
  • 团队对 RabbitMQ 更熟悉

最终选了 RabbitMQ。

踩过的坑

坑一:消息丢失

上线第二天发现,有些消息"消失"了。

排查原因:消费者处理完消息忘记手动确认,RabbitMQ 以为消费者还在处理,就重新投递。但消费者又是无状态部署,重启后就丢了。

解决

# 消费者处理完成后必须确认
def process_message(ch, method, properties, body):
    try:
        # 处理消息
        process(body)
        # 确认消息
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except Exception as e:
        # 拒绝消息,重新入队
        ch.basic_reject(delivery_tag=method.delivery_tag, requeue=True)

坑二:消息积压

某天下游服务响应变慢,消息队列积压了几十万条。

问题:消费者处理速度跟不上生产速度,消息越积越多。

解决

  • 增加消费者实例(但注意消费幂等性)
  • 优化下游服务性能
  • 考虑消息过期机制,超时丢弃旧消息

坑三:重复消费

由于网络问题或消费者重启,消息可能被重复投递。

解决:消费逻辑必须是幂等的。

def process_payment(message):
    order_id = message['order_id']
    payment_id = f"{order_id}_{message['timestamp']}"

    # 检查是否已处理
    if Payment.exists(payment_id):
        return

    # 处理支付
    payment = Payment.create(id=payment_id, ...)
    payment.save()

什么时候不该用消息队列

消息队列不是银弹,有些场景用了反而更复杂:

简单场景:服务调用链路简单,没必要引入复杂度

强一致性要求:金融交易等场景,要求数据实时一致

团队能力不足:如果团队没有运维经验,维护消息队列成本高

低延迟要求:如果要求实时响应,异步处理可能不合适

写在最后

消息队列这东西,解决了分布式系统的很多问题,但也引入了新的复杂度。

理解它的原理和最佳实践,才能用好它。否则就是给自己挖坑。

关键点

  • 根据业务场景选择合适的消息模型和产品
  • 重视消息的可靠性和幂等性处理
  • 建立完善的监控告警机制
  • 不是所有场景都需要消息队列

这次项目从零开始搭建消息队列,中间踩了不少坑。但回头看,理解了消息队列,对分布式系统的设计思路也清晰了很多。

版权声明: 本文首发于 指尖魔法屋-消息队列架构踩坑记录https://blog.thinkmoon.cn/post/25-message-queue-practice-kafka-rabbitmq/) 转载或引用必须申明原指尖魔法屋来源及源地址!