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 后,注册服务可以在用户创建成功后发送一条消息,例如 user.registered。邮件服务作为消费者订阅该消息,并在后台完成邮件发送。此时注册接口只需要确认核心业务已经完成,后续通知逻辑可以由消费者异步处理。
这种方式可以减少主流程等待时间,也能把耗时任务从核心链路中拆分出来。异步处理需要配套考虑任务失败、重试、幂等性和最终一致性。
系统解耦
在复杂业务系统中,一个业务事件往往会触发多个后续操作。例如订单创建后,可能涉及库存扣减、优惠券核销、积分更新、消息通知和风控检查。如果订单服务直接依赖这些下游服务,系统之间会形成较强的调用关系。
RabbitMQ 可以让服务之间通过消息间接协作。订单服务只负责发布“订单已创建”的消息,而不需要关心有多少个服务会消费该消息,也不需要知道这些服务的具体实现。库存服务、积分服务和通知服务可以分别订阅并处理自己关心的消息。
这种方式可以降低服务之间的直接依赖。新增、下线或调整消费者时,生产者通常不需要修改。
系统解耦不代表业务关系消失。消息格式、事件语义、消费顺序和异常处理仍然需要清晰定义。
流量削峰
RabbitMQ 也常用于流量削峰。秒杀活动、批量导入、集中推送或定时任务触发时,后端服务和数据库可能承受较高瞬时压力。
引入 RabbitMQ 后,请求可以先被转换为消息写入队列。后端消费者按照自身处理能力逐步消费消息,而不是被瞬时流量直接压垮。队列在这里起到缓冲层的作用,使系统可以用相对稳定的速度处理高峰期积压的任务。
例如在秒杀场景中,请求服务可以快速校验基本条件,然后将符合条件的请求写入 RabbitMQ。库存扣减和订单创建逻辑由后端消费者按顺序或按并发策略处理。这样可以减少数据库在瞬间承受的写入压力。
流量削峰不会提升系统总处理能力,只是把瞬时压力平滑到更长的时间窗口中。如果消息积压持续增长,仍然需要扩容、优化消费逻辑或限制入口流量。
事件通知与消息分发
RabbitMQ 支持将一条消息分发给多个消费者,因此也适合用于事件通知和消息广播。例如用户资料更新后,多个系统可能都需要感知这个事件:缓存系统需要刷新用户信息,搜索系统需要更新索引,审计系统需要记录变更日志。
此时生产者可以发布领域事件,多个消费者分别执行自己的处理逻辑。
这类场景通常会结合 RabbitMQ 的交换机和绑定机制实现。不同类型的交换机可以支持不同的消息分发方式,例如广播给所有队列、按路由键精确匹配,或按主题模式匹配多个消费者。
微服务之间的消息通信
在微服务架构中,服务之间既可能使用 HTTP、gRPC 等同步通信方式,也可能使用 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 重启后保留队列定义。队列持久化不等于消息持久化;消息本身也需要设置为持久化,并结合发布确认等机制。
Exchange:交换机
Exchange 负责接收生产者消息并进行路由。生产者通常先把消息发送到交换机,交换机再根据类型、路由键和绑定关系,将消息转发到一个或多个队列。
RabbitMQ 常见的交换机类型包括:
direct:根据路由键精确匹配队列fanout:将消息广播到所有绑定的队列topic:根据通配符规则匹配路由键headers:根据消息头属性进行匹配
交换机本身不负责长期保存消息。消息发送到交换机后如果没有匹配到任何队列,默认情况下可能会被丢弃。
交换机让生产者不需要直接依赖具体队列。
Binding:绑定关系
Binding 是交换机和队列之间的绑定关系。交换机通过绑定关系判断消息应该投递到哪些队列。
一个交换机可以绑定多个队列,一个队列也可以绑定到多个交换机。绑定时通常会指定一个绑定键。对于 direct 和 topic 类型的交换机,绑定键会参与路由匹配。

例如,一个 direct 交换机绑定了两个队列:
email_queue使用绑定键notification.emailsms_queue使用绑定键notification.sms
当生产者向交换机发送路由键为 notification.email 的消息时,交换机会把消息投递到 email_queue。当路由键为 notification.sms 时,消息会被投递到 sms_queue。
只要交换机和路由规则保持一致,系统可以通过新增或调整绑定来改变消费者接收消息的方式。
Routing Key:路由键
Routing Key 是生产者发送消息时携带的路由标识。交换机会根据路由键和绑定关系判断目标队列。
不同交换机类型对路由键的使用方式不同:
direct交换机会要求路由键和绑定键完全匹配topic交换机会根据通配符规则匹配路由键fanout交换机会忽略路由键,直接广播消息headers交换机主要根据消息头匹配,通常不依赖路由键
路由键通常采用层级化命名,例如 order.created、order.paid、user.registered 或 notification.email。这也方便在 topic 模式下进行规则匹配。
Message:消息
Message 是 RabbitMQ 中传递的数据单元。消息一般由消息体和消息属性组成。
消息体通常是业务数据,例如 JSON 字符串、文本内容或二进制数据。消息属性则用于描述消息的元信息,例如:
ContentType:消息内容类型,例如application/jsonDeliveryMode:消息是否持久化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 才会删除对应消息。处理失败时,消费者可以拒绝消息、重新入队,或将消息转入死信队列。

常见确认操作包括:
ack:确认消息已经成功处理nack:拒绝消息,可选择是否重新入队reject:拒绝单条消息,可选择是否重新入队
消息确认不能保证业务处理一定成功,但可以让系统在消费者异常、网络中断或进程崩溃时具备恢复能力。
RabbitMQ 的消息路由模型
RabbitMQ 的消息路由由交换机类型、路由键和绑定关系共同决定。生产者将消息发送到交换机,交换机再根据类型和绑定规则,将消息投递到一个或多个队列。
不同交换机类型对应不同的消息分发方式。

Direct Exchange
Direct Exchange 使用精确匹配的方式路由消息。生产者发送消息时会携带一个 routing key,队列绑定到交换机时也会指定一个 binding key。当 routing key 和 binding key 完全一致时,消息才会被投递到对应队列。
例如,一个 direct 交换机绑定了三个队列:
| 队列 | Binding Key |
|---|---|
order_created_queue | order.created |
order_paid_queue | order.paid |
order_canceled_queue | order.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.created、order.paid.success 或 user.profile.updated。
队列绑定到 topic 交换机时,可以使用通配符:
*:匹配一个单词#:匹配零个或多个单词
例如,一个 topic 交换机中存在以下绑定关系:
| 队列 | Binding Key | 含义 |
|---|---|---|
all_order_queue | order.# | 接收所有订单相关消息 |
order_paid_queue | order.paid.* | 接收订单支付相关消息 |
user_event_queue | user.* | 接收两段式用户事件 |
当 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=json 且 source=payment,并设置 x-match=all,那么只有同时满足这两个条件的消息才会进入该队列。如果设置为 x-match=any,只要其中一个条件匹配即可。
Headers Exchange 适合路由条件不方便表达为 routing key 的场景,例如根据消息来源、数据格式、租户标识、优先级或自定义属性进行分发。
headers 模型在普通业务系统中使用频率较低。多数场景可以通过 direct 或 topic 完成路由。
四种路由模型的对比
不同交换机类型没有绝对优劣,关键在于业务消息的分发方式。
| 交换机类型 | 路由方式 | 适合场景 |
|---|---|---|
direct | routing key 精确匹配 | 明确分类的任务或事件 |
fanout | 广播到所有绑定队列 | 事件广播、日志分发、缓存刷新 |
topic | routing 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-management5672 是 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;处理失败时调用 Nack 或 Reject。
完整消费者示例:
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 之前所有未确认消息。
处理消费失败
如果消费者处理消息失败,可以根据业务策略选择 Nack 或 Reject。
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_id 或 order_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 时,应把消息投递、消费确认、失败重试、死信处理、监控告警和消息结构作为一个整体设计。只有客户端代码接入完成,并不代表消息系统已经具备生产可用性。
