RabbitMQ 入门:使用场景、核心概念介绍
2026年06月27日

RabbitMQ 入门:使用场景、核心概念介绍

本文介绍了 RabbitMQ 的定位、常见使用场景、核心概念以及如何使用 Go 客户端进行消息发送与消费。

RabbitMQ 是什么

RabbitMQ 是一个开源消息代理,也常被归类为消息队列系统。它负责在应用、服务或系统模块之间传递消息,使发送方和接收方不必直接依赖彼此。

在一个典型系统中,某个服务完成业务操作后,可能需要通知其他服务继续处理后续任务。例如,订单服务创建订单后,库存服务需要扣减库存,通知服务需要发送消息,积分服务可能还需要更新用户积分。如果这些操作都由订单服务同步调用完成,系统之间的依赖会变得较强。一旦某个下游服务响应较慢或暂时不可用,主流程也容易受到影响。

RabbitMQ 提供基于消息的通信方式。发送方将消息发送到 RabbitMQ,RabbitMQ 按配置投递给对应接收方。接收方再消费消息并执行自身业务逻辑。这样可以削弱系统之间的直接调用关系。

从角色上看,RabbitMQ 更准确地说是 Message Broker,即消息代理。它不负责具体业务逻辑,而是负责接收、存储、路由和投递消息。生产者将消息发送给 RabbitMQ,消费者从 RabbitMQ 获取消息并处理。

RabbitMQ 最初实现的是 AMQP 协议,即 Advanced Message Queuing Protocol。AMQP 定义了一套消息传递模型,包括交换机、队列、绑定、路由键、消息确认等概念。因此,RabbitMQ 不只是先进先出队列,而是支持多种路由模式和可靠性控制的消息系统。

在实际项目中,RabbitMQ 常用于异步处理、服务解耦、流量削峰、事件通知和任务分发等场景。它适合处理业务系统中需要可靠传递、可控消费和灵活路由的消息。对于需要极高吞吐量、长期日志存储或流式处理的场景,系统通常还会评估 Kafka、Pulsar 等其他消息系统。

RabbitMQ 的核心价值在于提供消息中转层,使不同服务通过消息协作,而不是彼此强绑定地同步调用。

RabbitMQ 整体架构

RabbitMQ 的典型使用场景

RabbitMQ 通常用于系统间消息传递、异步任务和服务解耦。它不直接改变业务规则,而是改变模块之间的协作方式。

异步任务处理

异步任务处理是 RabbitMQ 最常见的使用场景之一。发送邮件、发送短信、生成报表、同步搜索索引、处理图片或推送通知等操作,通常不需要阻塞主流程。

以用户注册为例,系统在创建用户账号后,可能需要发送欢迎邮件。如果注册接口直接调用邮件服务,接口响应时间会受到邮件服务处理速度的影响。当邮件服务出现短暂故障时,注册流程也可能被拖慢甚至失败。

使用 RabbitMQ 后,注册服务可以在用户创建成功后发送一条消息,例如 user.registered。邮件服务作为消费者订阅该消息,并在后台完成邮件发送。此时注册接口只需要确认核心业务已经完成,后续通知逻辑可以由消费者异步处理。

这种方式可以减少主流程等待时间,也能把耗时任务从核心链路中拆分出来。异步处理需要配套考虑任务失败、重试、幂等性和最终一致性。

系统解耦

在复杂业务系统中,一个业务事件往往会触发多个后续操作。例如订单创建后,可能涉及库存扣减、优惠券核销、积分更新、消息通知和风控检查。如果订单服务直接依赖这些下游服务,系统之间会形成较强的调用关系。

RabbitMQ 可以让服务之间通过消息间接协作。订单服务只负责发布“订单已创建”的消息,而不需要关心有多少个服务会消费该消息,也不需要知道这些服务的具体实现。库存服务、积分服务和通知服务可以分别订阅并处理自己关心的消息。

这种方式可以降低服务之间的直接依赖。新增、下线或调整消费者时,生产者通常不需要修改。

系统解耦不代表业务关系消失。消息格式、事件语义、消费顺序和异常处理仍然需要清晰定义。

流量削峰

RabbitMQ 也常用于流量削峰。秒杀活动、批量导入、集中推送或定时任务触发时,后端服务和数据库可能承受较高瞬时压力。

引入 RabbitMQ 后,请求可以先被转换为消息写入队列。后端消费者按照自身处理能力逐步消费消息,而不是被瞬时流量直接压垮。队列在这里起到缓冲层的作用,使系统可以用相对稳定的速度处理高峰期积压的任务。

例如在秒杀场景中,请求服务可以快速校验基本条件,然后将符合条件的请求写入 RabbitMQ。库存扣减和订单创建逻辑由后端消费者按顺序或按并发策略处理。这样可以减少数据库在瞬间承受的写入压力。

流量削峰不会提升系统总处理能力,只是把瞬时压力平滑到更长的时间窗口中。如果消息积压持续增长,仍然需要扩容、优化消费逻辑或限制入口流量。

事件通知与消息分发

RabbitMQ 支持将一条消息分发给多个消费者,因此也适合用于事件通知和消息广播。例如用户资料更新后,多个系统可能都需要感知这个事件:缓存系统需要刷新用户信息,搜索系统需要更新索引,审计系统需要记录变更日志。

此时生产者可以发布领域事件,多个消费者分别执行自己的处理逻辑。

这类场景通常会结合 RabbitMQ 的交换机和绑定机制实现。不同类型的交换机可以支持不同的消息分发方式,例如广播给所有队列、按路由键精确匹配,或按主题模式匹配多个消费者。

微服务之间的消息通信

在微服务架构中,服务之间既可能使用 HTTP、gRPC 等同步通信方式,也可能使用 RabbitMQ 这样的消息系统进行异步通信。RabbitMQ 适合用于不要求立即返回结果、但需要可靠传递业务事件的场景。

例如支付服务在支付成功后,可以发布支付成功事件。订单服务消费该事件后更新订单状态,通知服务消费该事件后发送支付成功通知,财务系统消费该事件后生成对账记录。各服务通过消息进行协作,而不是由支付服务逐个调用下游接口。

这种通信方式可以提升系统弹性。消费者暂时不可用时,只要消息被 RabbitMQ 保存,消费者恢复后仍可继续处理。微服务消息通信需要设计清晰的事件模型,并处理重复消息、消息乱序、消费失败和数据一致性。

任务分发与工作队列

RabbitMQ 也可以作为工作队列使用。生产者持续向队列中投递任务,多个消费者共同从队列中取出任务并处理。RabbitMQ 会将消息分发给不同消费者,使任务可以被并行处理。

这种模式适合处理批量任务、后台计算、文件转换、数据同步等工作。例如系统需要处理一批图片压缩任务,可以将每张图片的处理请求写入队列,再启动多个消费者并发执行压缩逻辑。

在工作队列模式中,消费者数量可以根据任务积压情况调整。配合消息确认机制,RabbitMQ 可以在消费者异常退出时重新投递未确认的消息。

RabbitMQ 适合业务系统中的异步协作、事件分发和任务缓冲。它的重点不只是“把消息放进队列”,而是通过消息代理降低系统耦合。

RabbitMQ 的核心概念

RabbitMQ 的常见消息流程是:生产者将消息发送到交换机,交换机根据规则将消息路由到一个或多个队列,消费者再从队列中取出消息并处理。

生产者负责发送消息,交换机负责路由消息,队列负责保存消息,消费者负责处理消息。绑定关系和路由键决定消息会被投递到哪些队列。

RabbitMQ 消息生命周期

Producer:生产者

Producer 指消息发送方,即向 RabbitMQ 投递消息的应用程序。生产者通常是业务系统中的某个模块或服务,例如订单服务、用户服务、支付服务或后台任务服务。

生产者的职责通常包括:

  • 创建消息内容
  • 指定消息要发送到哪个交换机
  • 设置路由键或消息属性
  • 根据需要设置消息是否持久化
  • 处理发送失败或连接异常

生产者不直接关心最终由哪个消费者处理消息。消息进入哪个队列、被哪些消费者消费,由交换机、绑定关系和消费者订阅情况决定。

例如,订单服务可以在订单创建成功后发送一条 order.created 消息,而不需要逐个调用通知服务、积分服务或库存服务。

Consumer:消费者

Consumer 指消息接收方,即从 RabbitMQ 队列中获取消息并执行处理逻辑的应用程序。消费者可以是长期运行的服务,也可以是后台任务进程。

消费者订阅队列并监听消息。RabbitMQ 将消息投递给消费者后,消费者根据消息内容执行具体业务逻辑,例如发送邮件、更新订单状态、写入日志或同步数据。

在可靠性要求较高的场景中,消费者通常在业务逻辑执行成功后手动确认消息,即 ACK。RabbitMQ 收到确认后,才会认为消息已经被成功处理。

如果消费者在处理过程中异常退出,且尚未确认消息,RabbitMQ 可以重新投递该消息。

Queue:队列

Queue 是 RabbitMQ 中用于存储消息的结构。消息在被消费者处理之前,通常会先进入队列等待消费。

队列可以理解为生产者和消费者之间的缓冲区。生产者发送消息的速度和消费者处理消息的速度不一定一致。当生产者发送速度更快时,消息会在队列中积压;当消费者处理能力足够时,队列中的消息会逐步减少。

RabbitMQ Queue 队列

队列具有一些重要属性,例如:

  • 是否持久化
  • 是否排他
  • 是否自动删除
  • 是否设置最大长度
  • 是否绑定死信交换机

持久化队列可以在 RabbitMQ 重启后保留队列定义。队列持久化不等于消息持久化;消息本身也需要设置为持久化,并结合发布确认等机制。

Exchange:交换机

Exchange 负责接收生产者消息并进行路由。生产者通常先把消息发送到交换机,交换机再根据类型、路由键和绑定关系,将消息转发到一个或多个队列。

RabbitMQ 常见的交换机类型包括:

  • direct:根据路由键精确匹配队列
  • fanout:将消息广播到所有绑定的队列
  • topic:根据通配符规则匹配路由键
  • headers:根据消息头属性进行匹配

交换机本身不负责长期保存消息。消息发送到交换机后如果没有匹配到任何队列,默认情况下可能会被丢弃。

交换机让生产者不需要直接依赖具体队列。

Binding:绑定关系

Binding 是交换机和队列之间的绑定关系。交换机通过绑定关系判断消息应该投递到哪些队列。

一个交换机可以绑定多个队列,一个队列也可以绑定到多个交换机。绑定时通常会指定一个绑定键。对于 directtopic 类型的交换机,绑定键会参与路由匹配。

RabbitMQ Binding 绑定关系

例如,一个 direct 交换机绑定了两个队列:

  • email_queue 使用绑定键 notification.email
  • sms_queue 使用绑定键 notification.sms

当生产者向交换机发送路由键为 notification.email 的消息时,交换机会把消息投递到 email_queue。当路由键为 notification.sms 时,消息会被投递到 sms_queue

只要交换机和路由规则保持一致,系统可以通过新增或调整绑定来改变消费者接收消息的方式。

Routing Key:路由键

Routing Key 是生产者发送消息时携带的路由标识。交换机会根据路由键和绑定关系判断目标队列。

不同交换机类型对路由键的使用方式不同:

  • direct 交换机会要求路由键和绑定键完全匹配
  • topic 交换机会根据通配符规则匹配路由键
  • fanout 交换机会忽略路由键,直接广播消息
  • headers 交换机主要根据消息头匹配,通常不依赖路由键

路由键通常采用层级化命名,例如 order.createdorder.paiduser.registerednotification.email。这也方便在 topic 模式下进行规则匹配。

Message:消息

Message 是 RabbitMQ 中传递的数据单元。消息一般由消息体和消息属性组成。

消息体通常是业务数据,例如 JSON 字符串、文本内容或二进制数据。消息属性则用于描述消息的元信息,例如:

  • ContentType:消息内容类型,例如 application/json
  • DeliveryMode:消息是否持久化
  • Timestamp:消息创建时间
  • CorrelationId:关联请求或链路追踪标识
  • ReplyTo:用于请求响应模式的回复队列
  • Headers:自定义消息头

RabbitMQ 不理解消息体的业务含义。消息格式、字段含义和兼容性需要由生产者和消费者共同约定。

跨服务消息需要考虑版本兼容、字段扩展、序列化格式和敏感信息处理。

Virtual Host:虚拟主机

Virtual Host,简称 vhost,是 RabbitMQ 中的逻辑隔离单元。一个 RabbitMQ 实例可以创建多个 vhost,每个 vhost 有独立的交换机、队列、绑定和权限配置。

vhost 可以用于区分不同环境、不同业务系统或不同租户。例如,开发环境、测试环境和生产环境可以使用不同的 vhost;多个业务团队也可以在同一个 RabbitMQ 集群中使用不同 vhost 实现隔离。

权限控制通常也基于 vhost 配置。RabbitMQ 用户可以被授予某个 vhost 下的读、写和配置权限。

默认情况下,RabbitMQ 会提供名为 / 的 vhost。多系统或多环境部署通常会显式创建独立 vhost。

Channel:信道

Channel 是建立在 RabbitMQ 连接之上的轻量级通信通道。客户端通常先创建 TCP 连接,再在连接上创建一个或多个 channel。声明队列、声明交换机、发布消息和消费消息等 AMQP 操作都通过 channel 完成。

TCP 连接相对较重。通过在一个连接上复用多个 channel,客户端可以更高效地与 RabbitMQ 通信。

在 Go 客户端中,应用通常先建立连接,再打开 channel,然后通过 channel 声明队列、发布消息和消费消息。channel 不是线程安全的;多个 goroutine 并发发布时,需要额外同步或创建独立 channel。

消息确认机制

消息确认机制用于让 RabbitMQ 判断消息是否被消费者成功处理。自动确认模式下,RabbitMQ 只要把消息投递给消费者,就认为消息已经处理完成;如果消费者收到消息后崩溃,消息可能丢失。

可靠消费通常使用手动确认。消费者处理成功后调用 ACK,RabbitMQ 才会删除对应消息。处理失败时,消费者可以拒绝消息、重新入队,或将消息转入死信队列。

RabbitMQ 消息确认 ACK

常见确认操作包括:

  • ack:确认消息已经成功处理
  • nack:拒绝消息,可选择是否重新入队
  • reject:拒绝单条消息,可选择是否重新入队

消息确认不能保证业务处理一定成功,但可以让系统在消费者异常、网络中断或进程崩溃时具备恢复能力。

RabbitMQ 的消息路由模型

RabbitMQ 的消息路由由交换机类型、路由键和绑定关系共同决定。生产者将消息发送到交换机,交换机再根据类型和绑定规则,将消息投递到一个或多个队列。

不同交换机类型对应不同的消息分发方式。

RabbitMQ Exchange 类型

Direct Exchange

Direct Exchange 使用精确匹配的方式路由消息。生产者发送消息时会携带一个 routing key,队列绑定到交换机时也会指定一个 binding key。当 routing key 和 binding key 完全一致时,消息才会被投递到对应队列。

例如,一个 direct 交换机绑定了三个队列:

队列Binding Key
order_created_queueorder.created
order_paid_queueorder.paid
order_canceled_queueorder.canceled

当生产者发送 routing key 为 order.paid 的消息时,消息只会进入绑定键为 order.paid 的队列。其他队列不会收到这条消息。

Direct Exchange 适合路由规则明确、消息类型可以被精确分类的场景,例如订单状态变更、通知渠道区分、任务类型分发等。

direct 模型的灵活性相对有限。如果消费者需要订阅一类具有共同前缀的消息,例如所有订单相关事件,topic 模型通常更合适。

Fanout Exchange

Fanout Exchange 使用广播方式路由消息。它会忽略 routing key,将消息投递到所有绑定到该交换机的队列。

例如,一个 fanout 交换机绑定了三个队列:

队列用途
cache_refresh_queue刷新缓存
search_index_queue更新搜索索引
audit_log_queue记录审计日志

当生产者向这个交换机发送一条用户资料更新消息时,三个队列都会收到同一条消息。每个消费者可以根据自己的职责执行对应逻辑。

Fanout Exchange 适合事件广播场景,例如配置变更通知、缓存刷新、系统公告、日志分发或领域事件广播。

它不适合精细筛选消息。只要队列绑定到了 fanout 交换机,就会收到所有消息。

Topic Exchange

Topic Exchange 使用通配符匹配的方式路由消息。它仍然依赖 routing key,但 routing key 通常会按照层级结构命名,并使用点号分隔,例如 order.createdorder.paid.successuser.profile.updated

队列绑定到 topic 交换机时,可以使用通配符:

  • *:匹配一个单词
  • #:匹配零个或多个单词

例如,一个 topic 交换机中存在以下绑定关系:

队列Binding Key含义
all_order_queueorder.#接收所有订单相关消息
order_paid_queueorder.paid.*接收订单支付相关消息
user_event_queueuser.*接收两段式用户事件

当 routing key 为 order.paid.success 时,它可以匹配 order.#,也可以匹配 order.paid.*。因此,同一条消息可能被投递到多个队列。

Topic Exchange 适合消息类型较多、路由规则具有层级结构的场景。例如电商系统中的订单事件、支付事件、用户事件和库存事件,都可以通过统一命名规则组织。

topic 模型比 direct 更灵活。消费者可以订阅某个具体事件,也可以订阅某类事件。代价是 routing key 命名必须保持一致。

Headers Exchange

Headers Exchange 根据消息头属性进行路由,而不是主要依赖 routing key。生产者发送消息时可以设置 headers,队列绑定时也可以指定需要匹配的 headers 条件。

Headers Exchange 通常支持两种匹配方式:

  • x-match: all:所有指定 header 都匹配时才投递
  • x-match: any:任意一个指定 header 匹配时即可投递

例如,一条消息包含以下 headers:

{
  "format": "json",
  "source": "payment",
  "priority": "high"
}

如果某个队列绑定时要求 format=jsonsource=payment,并设置 x-match=all,那么只有同时满足这两个条件的消息才会进入该队列。如果设置为 x-match=any,只要其中一个条件匹配即可。

Headers Exchange 适合路由条件不方便表达为 routing key 的场景,例如根据消息来源、数据格式、租户标识、优先级或自定义属性进行分发。

headers 模型在普通业务系统中使用频率较低。多数场景可以通过 direct 或 topic 完成路由。

四种路由模型的对比

不同交换机类型没有绝对优劣,关键在于业务消息的分发方式。

交换机类型路由方式适合场景
directrouting key 精确匹配明确分类的任务或事件
fanout广播到所有绑定队列事件广播、日志分发、缓存刷新
topicrouting key 通配符匹配层级化事件、复杂订阅关系
headers根据消息头属性匹配多属性条件路由

direct 和 topic 是常见选择。direct 适合规则简单且稳定的消息分发,topic 适合消息类型较多、需要按层级订阅的场景。fanout 适合广播通知,headers 适合基于消息属性组合进行路由的场景。

路由模型选择建议

如果系统只是希望把某类任务发送到固定队列,可以优先考虑 direct 交换机。

如果系统希望一个事件被多个服务同时感知,可以考虑 fanout 交换机。

如果系统中存在大量领域事件,并且消费者需要按事件类别订阅,可以考虑 topic 交换机。

如果消息路由依赖多个属性组合,而不是单个 routing key,可以考虑 headers 交换机。

RabbitMQ 的路由能力来自交换机和绑定关系的组合。生产者发布消息,消费者通过队列接收消息,交换机根据规则完成分发。

使用 Go 操作 RabbitMQ

安装 Go 客户端

go get github.com/rabbitmq/amqp091-go

准备 RabbitMQ 服务

可以使用 Docker 启动一个带管理后台的 RabbitMQ 容器:

docker run --name rabbitmq-demo \
  -p 5672:5672 \
  -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=guest \
  -e RABBITMQ_DEFAULT_PASS=guest \
  rabbitmq:3-management

5672 是 AMQP 协议端口,15672 是管理后台端口。管理后台地址为 http://localhost:15672,默认用户名和密码都是 guest

建立连接和 Channel

Go 程序通常先创建连接,再基于连接打开 channel。后续声明队列、发布消息和消费消息都通过 channel 完成。

package main

import (
    "log"

    amqp "github.com/rabbitmq/amqp091-go"
)

func main() {
    conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    if err != nil {
        log.Fatalf("connect rabbitmq: %v", err)
    }
    defer conn.Close()

    ch, err := conn.Channel()
    if err != nil {
        log.Fatalf("open channel: %v", err)
    }
    defer ch.Close()

    log.Println("connected to RabbitMQ")
}

连接地址格式:

amqp://username:password@host:port/vhost

默认 vhost / 可以直接写成连接地址末尾的 /。自定义 vhost 需要写入对应名称,并注意 URL 编码。

声明队列

生产者发送消息前,通常会先声明队列。队列声明是幂等的:队列不存在时会创建;队列已存在且参数一致时会复用;参数不一致时会返回错误。

q, err := ch.QueueDeclare(
    "hello_queue", // name
    true,          // durable
    false,         // delete when unused
    false,         // exclusive
    false,         // no-wait
    nil,           // arguments
)
if err != nil {
    log.Fatalf("declare queue: %v", err)
}

log.Printf("queue declared: %s", q.Name)

durable 设置为 true 表示队列定义会被持久化。队列持久化不等于消息持久化,消息也需要设置 DeliveryMode

发送消息

生产者可以使用默认交换机发送消息。RabbitMQ 的默认交换机名称是空字符串 ""。使用默认交换机时,routing key 会被当作队列名称,消息会被投递到同名队列。

err = ch.Publish(
    "",           // exchange
    q.Name,       // routing key
    false,        // mandatory
    false,        // immediate
    amqp.Publishing{
        ContentType:  "text/plain",
        DeliveryMode: amqp.Persistent,
        Body:         []byte("hello rabbitmq"),
    },
)
if err != nil {
    log.Fatalf("publish message: %v", err)
}

log.Println("message published")

DeliveryMode: amqp.Persistent 表示消息以持久化方式发布。完整可靠性还依赖发布确认、磁盘刷写策略和集群配置。

完整生产者示例:

package main

import (
    "log"

    amqp "github.com/rabbitmq/amqp091-go"
)

func main() {
    conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    if err != nil {
        log.Fatalf("connect rabbitmq: %v", err)
    }
    defer conn.Close()

    ch, err := conn.Channel()
    if err != nil {
        log.Fatalf("open channel: %v", err)
    }
    defer ch.Close()

    q, err := ch.QueueDeclare(
        "hello_queue",
        true,
        false,
        false,
        false,
        nil,
    )
    if err != nil {
        log.Fatalf("declare queue: %v", err)
    }

    err = ch.Publish(
        "",
        q.Name,
        false,
        false,
        amqp.Publishing{
            ContentType:  "text/plain",
            DeliveryMode: amqp.Persistent,
            Body:         []byte("hello rabbitmq"),
        },
    )
    if err != nil {
        log.Fatalf("publish message: %v", err)
    }

    log.Println("message published")
}

消费消息

消费者通过 Consume 方法订阅队列中的消息。

msgs, err := ch.Consume(
    q.Name, // queue
    "",     // consumer
    false,  // auto-ack
    false,  // exclusive
    false,  // no-local
    false,  // no-wait
    nil,    // args
)
if err != nil {
    log.Fatalf("consume message: %v", err)
}

auto-ack 设置为 false 表示使用手动确认。处理成功后调用 Ack;处理失败时调用 NackReject

完整消费者示例:

package main

import (
    "log"

    amqp "github.com/rabbitmq/amqp091-go"
)

func main() {
    conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    if err != nil {
        log.Fatalf("connect rabbitmq: %v", err)
    }
    defer conn.Close()

    ch, err := conn.Channel()
    if err != nil {
        log.Fatalf("open channel: %v", err)
    }
    defer ch.Close()

    q, err := ch.QueueDeclare(
        "hello_queue",
        true,
        false,
        false,
        false,
        nil,
    )
    if err != nil {
        log.Fatalf("declare queue: %v", err)
    }

    msgs, err := ch.Consume(
        q.Name,
        "",
        false,
        false,
        false,
        false,
        nil,
    )
    if err != nil {
        log.Fatalf("consume message: %v", err)
    }

    log.Println("waiting for messages")

    forever := make(chan struct{})

    go func() {
        for msg := range msgs {
            log.Printf("received: %s", string(msg.Body))

            if err := msg.Ack(false); err != nil {
                log.Printf("ack message: %v", err)
            }
        }
    }()

    <-forever
}

Ack(false) 表示只确认当前消息。参数为 true 时,会确认当前 delivery tag 之前所有未确认消息。

处理消费失败

如果消费者处理消息失败,可以根据业务策略选择 NackReject

for msg := range msgs {
    if err := handleMessage(msg.Body); err != nil {
        log.Printf("handle message failed: %v", err)

        if nackErr := msg.Nack(false, true); nackErr != nil {
            log.Printf("nack message: %v", nackErr)
        }
        continue
    }

    if err := msg.Ack(false); err != nil {
        log.Printf("ack message: %v", err)
    }
}

msg.Nack(false, true) 的第二个参数表示是否重新入队。设置为 true 时,消息会回到队列;设置为 false 时,如果队列配置了死信交换机,消息可以进入死信队列。

不建议对同一条失败消息无限重新入队。常见做法是结合重试次数、延迟队列或死信队列处理失败消息。

声明交换机并绑定队列

实际项目中,通常会显式声明交换机,再将队列绑定到交换机。以下代码声明一个 direct 交换机,并将队列绑定到 order.created 路由键:

err = ch.ExchangeDeclare(
    "order_exchange", // name
    "direct",         // type
    true,             // durable
    false,            // auto-deleted
    false,            // internal
    false,            // no-wait
    nil,              // arguments
)
if err != nil {
    log.Fatalf("declare exchange: %v", err)
}

q, err := ch.QueueDeclare(
    "order_created_queue",
    true,
    false,
    false,
    false,
    nil,
)
if err != nil {
    log.Fatalf("declare queue: %v", err)
}

err = ch.QueueBind(
    q.Name,           // queue name
    "order.created", // routing key
    "order_exchange", // exchange
    false,
    nil,
)
if err != nil {
    log.Fatalf("bind queue: %v", err)
}

生产者发送消息时指定交换机和 routing key:

err = ch.Publish(
    "order_exchange",
    "order.created",
    false,
    false,
    amqp.Publishing{
        ContentType:  "application/json",
        DeliveryMode: amqp.Persistent,
        Body: []byte(`{
            "order_id": "202606270001",
            "user_id": "10001",
            "amount": 19900
        }`),
    },
)
if err != nil {
    log.Fatalf("publish order event: %v", err)
}

生产者发布订单创建事件,交换机根据 order.created 路由键将消息投递到绑定的订单队列。其他服务需要处理该事件时,可以新增队列并绑定到同一个交换机,不必修改生产者。

设置消费预取数量

为了避免某个消费者一次拿走过多消息,可以设置 QoS,也就是预取数量。

err = ch.Qos(
    1,     // prefetch count
    0,     // prefetch size
    false, // global
)
if err != nil {
    log.Fatalf("set qos: %v", err)
}

prefetch count 设置为 1 时,RabbitMQ 在消费者确认当前消息之前,不会继续向该消费者投递新消息。

预取数量需要结合消息处理耗时、消费者数量、内存占用和吞吐量调整。

使用中的基本建议

Go 程序接入 RabbitMQ 时的基本原则:

  • 连接和 channel 都需要在不再使用时关闭
  • 队列、交换机和绑定声明应保持参数一致
  • 可靠消费场景应优先使用手动 ACK
  • 生产者需要考虑发布失败和发布确认
  • 消费者需要处理业务失败、重复消息和幂等性
  • 长期运行的服务需要处理连接断开和自动重连
  • 多个 goroutine 并发发布消息时,不应随意共享同一个 channel

RabbitMQ 客户端代码本身并不复杂,真正需要谨慎设计的是消息语义、失败处理和服务之间的数据一致性。

RabbitMQ 使用中的注意事项

RabbitMQ 接入业务系统后,重点不在 API 调用本身,而在可靠性、可观测性和失败处理。以下问题通常需要在设计阶段明确。

消息持久化

消息持久化需要同时满足两个条件:

  • 队列设置为 durable
  • 消息设置为 persistent

Go 客户端中,队列声明时需要将 durable 设置为 true

q, err := ch.QueueDeclare(
    "task_queue",
    true,
    false,
    false,
    false,
    nil,
)
if err != nil {
    log.Fatalf("declare queue: %v", err)
}

发布消息时设置 DeliveryMode

err = ch.Publish(
    "",
    q.Name,
    false,
    false,
    amqp.Publishing{
        ContentType:  "application/json",
        DeliveryMode: amqp.Persistent,
        Body:         []byte(`{"task_id":"10001"}`),
    },
)
if err != nil {
    log.Fatalf("publish message: %v", err)
}

持久化可以降低 RabbitMQ 重启导致消息丢失的风险,但它不是完整的可靠投递方案。生产者还需要结合发布确认,消费者也需要使用手动 ACK。

发布确认

生产者调用 Publish 成功,只表示客户端已经把消息写入连接,并不等同于 RabbitMQ 已经可靠接收并处理该消息。对可靠性要求较高的场景,应启用 publisher confirm。

if err := ch.Confirm(false); err != nil {
    log.Fatalf("enable publisher confirms: %v", err)
}

confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 1))

err = ch.Publish(
    "",
    q.Name,
    false,
    false,
    amqp.Publishing{
        ContentType:  "text/plain",
        DeliveryMode: amqp.Persistent,
        Body:         []byte("important message"),
    },
)
if err != nil {
    log.Fatalf("publish message: %v", err)
}

confirm := <-confirms
if !confirm.Ack {
    log.Fatal("message was not acknowledged by broker")
}

发布确认会增加实现复杂度和延迟。批量发布时通常会结合异步确认、超时控制和重试策略。

手动 ACK

消费者建议关闭自动确认。自动确认模式下,RabbitMQ 投递消息后即认为消息处理完成。如果消费者在业务逻辑执行前崩溃,消息不会再次投递。

msgs, err := ch.Consume(
    q.Name,
    "",
    false,
    false,
    false,
    false,
    nil,
)
if err != nil {
    log.Fatalf("consume message: %v", err)
}

for msg := range msgs {
    if err := handleMessage(msg.Body); err != nil {
        _ = msg.Nack(false, false)
        continue
    }

    _ = msg.Ack(false)
}

ACK 的位置应该放在业务逻辑成功之后。对于写数据库、调用外部服务等操作,需要先确认业务侧已经完成,再确认消息。

幂等性

RabbitMQ 不能保证业务层“只处理一次”。网络中断、消费者崩溃、ACK 丢失、重试和重新入队都可能导致重复消费。

消费者处理逻辑应具备幂等性。常见做法包括:

  • 使用业务唯一键去重
  • 在数据库中记录消息处理状态
  • 对状态变更使用条件更新
  • 对外部调用引入幂等键

例如订单支付成功事件可以使用 payment_idorder_id + event_type 作为幂等键。消费者在处理前先检查该事件是否已经处理过,避免重复更新订单状态或重复发放权益。

重试机制

消费失败不应简单地无限重新入队。无限重试会让异常消息反复进入消费者,影响正常消息处理。

常见重试策略包括:

  • 立即重试:适合短暂、低成本失败
  • 延迟重试:适合依赖外部服务恢复的失败
  • 固定次数重试:超过次数后进入死信队列
  • 人工介入:适合数据异常或业务状态冲突

RabbitMQ 本身没有内置统一的重试语义,通常通过死信交换机、TTL 队列、延迟消息插件或业务侧重试表实现。

死信队列

死信队列用于接收无法正常消费的消息。消息进入死信队列的常见原因包括:

  • 消息被拒绝且不重新入队
  • 消息过期
  • 队列达到最大长度

声明业务队列时可以指定死信交换机:

args := amqp.Table{
    "x-dead-letter-exchange":    "dlx_exchange",
    "x-dead-letter-routing-key": "task.dead",
}

_, err := ch.QueueDeclare(
    "task_queue",
    true,
    false,
    false,
    false,
    args,
)
if err != nil {
    log.Fatalf("declare queue: %v", err)
}

死信队列不只是“失败消息归档”。它通常还承担问题排查、人工补偿、延迟重试入口等职责。

消费者并发

消费者并发可以通过增加消费者进程、增加 goroutine 或提高单个消费者的预取数量实现。并发策略需要结合消息处理耗时、消息顺序要求和下游系统承载能力。

如果消息之间没有顺序依赖,可以增加消费者数量提升吞吐量。如果同一业务实体需要按顺序处理,例如同一个订单的多个状态事件,就需要避免并发导致乱序更新。

Qos 可以限制单个消费者未确认消息数量:

err := ch.Qos(
    10,
    0,
    false,
)
if err != nil {
    log.Fatalf("set qos: %v", err)
}

预取数量过小会限制吞吐,过大可能导致消息分配不均、消费者内存升高或失败后重新投递大量消息。

连接管理

长期运行的服务需要处理连接断开、channel 异常和 RabbitMQ 节点重启。amqp091-go 不会自动重连,需要业务代码自行实现重连逻辑。

重连时需要重新完成以下操作:

  • 建立连接
  • 创建 channel
  • 声明交换机
  • 声明队列
  • 绑定队列
  • 恢复消费

生产者和消费者都不应假设连接永久可用。发布、消费、ACK、NACK 等操作都可能返回错误。

监控与告警

RabbitMQ 运行状态需要纳入监控。常见指标包括:

  • 队列消息积压量
  • 消费速率
  • 发布速率
  • 未确认消息数量
  • 消费者数量
  • 连接数和 channel 数
  • 内存、磁盘和文件描述符使用情况
  • 死信队列消息数量

队列积压持续增长通常意味着消费者处理能力不足、下游依赖变慢或出现异常消息。未确认消息数量持续增长,通常需要检查消费者处理耗时、ACK 逻辑和连接状态。

消息设计

消息结构应尽量稳定、明确,避免把内部数据库结构直接暴露给其他服务。跨服务事件更适合表达“发生了什么”,而不是要求消费者执行某个具体命令。

例如:

{
  "event_id": "evt_202606270001",
  "event_type": "order.created",
  "occurred_at": "2026-06-27T10:30:00Z",
  "data": {
    "order_id": "202606270001",
    "user_id": "10001",
    "amount": 19900
  }
}

event_id 可用于幂等处理,event_type 表达事件类型,occurred_at 记录事件发生时间,data 保存业务数据。后续字段扩展应保持向后兼容。

小结

RabbitMQ 适合承担业务系统中的异步协作、事件分发和任务缓冲。它能降低服务之间的直接耦合,但不会自动解决一致性、幂等性、重试和运维问题。

使用 RabbitMQ 时,应把消息投递、消费确认、失败重试、死信处理、监控告警和消息结构作为一个整体设计。只有客户端代码接入完成,并不代表消息系统已经具备生产可用性。