
在 go 中使用 rabbitmq fanout exchange 时,若多个消费者仅交替接收消息而非同时接收,通常是因为未正确定义交换器(exchange)类型或队列绑定逻辑错误;需显式声明 fanout 类型交换器并确保各消费者绑定到同一交换器下的独立队列。
在 go 中使用 rabbitmq fanout exchange 时,若多个消费者仅交替接收消息而非同时接收,通常是因为未正确定义交换器(exchange)类型或队列绑定逻辑错误;需显式声明 fanout 类型交换器并确保各消费者绑定到同一交换器下的独立队列。
Fanout Exchange 的核心语义是“广播”——所有绑定到该交换器的队列都会完整复制并收到每一条发布消息。但要实现这一行为,必须满足两个前提条件:
- 交换器必须被正确定义为 fanout 类型(而非默认的 direct 或未声明);
- 每个消费者应使用各自独立的队列名(而非共用 "example.queue"),并将其绑定到同一个 fanout 交换器。
您当前代码中的关键问题在于:
- ❌ 未调用 channel.ExchangeDeclare(...) 声明 fanout 类型交换器;
- ❌ 两个消费者均使用了相同的队列名 "example.queue",导致 RabbitMQ 将其视为同一个队列的两个消费者实例(即竞争消费模式),因此消息被轮询分发(Round-Robin),而非广播。
✅ 正确做法如下:
✅ 步骤一:统一声明 Fanout Exchange
在发布端(或消费者初始化时)提前声明交换器(建议在连接建立后、消费前执行一次即可):
RabbitMQ 4.2.3 是 2026 年初发布的重要稳定更新版本,重点修复了 Khepri 元数据存储相关问题,并改进了监控性能。对于使用 Docker、Kubernetes 或微服务架构的开发团队来说,该版本兼容性和稳定性表现较好。
err := channel.ExchangeDeclare(
"logs", // 交换器名称(推荐语义化命名,如 "logs"、"broadcast")
"fanout", // 类型必须为 "fanout"
true, // durable: 持久化,重启后仍存在
false, // auto-deleted: 不自动删除
false, // internal: 非内部交换器
false, // no-wait
nil, // arguments
)
if err != nil {
log.Fatalf("Failed to declare exchange: %v", err)
}
⚠️ 注意:ExchangeDeclare 只需调用一次(幂等),多个消费者/生产者可复用同一交换器。
✅ 步骤二:为每个消费者创建专属队列并绑定
修改您的 HandleMessageFanout1 和 HandleMessageFanout2,使用不同队列名,并显式绑定到 logs 交换器:
// HandleMessageFanout1 —— 使用队列 "queue-fanout-1"
func HandleMessageFanout1() {
conn := system.EltropyAppContext.RabbitMQConn
ch, err := conn.Channel()
if err != nil {
log.Fatalf("Failed to open channel: %v", err)
}
defer ch.Close()
// 声明专属队列(不指定名称则由 RabbitMQ 自动生成,此处显式命名便于调试)
q, err := ch.QueueDeclare(
"queue-fanout-1", // 唯一队列名
true, // durable
false, // delete when unused
false, // exclusive
false, // no-wait
nil, // args
)
if err != nil {
log.Fatalf("Failed to declare queue: %v", err)
}
// 绑定队列到 fanout 交换器(routingKey 在 fanout 中被忽略,传空字符串即可)
err = ch.QueueBind(
q.Name, // queue name
"", // routing key (ignored for fanout)
"logs", // exchange name
false, // no-wait
nil,
)
if err != nil {
log.Fatalf("Failed to bind queue to exchange: %v", err)
}
// 开始消费
msgs, err := ch.Consume(
q.Name, // queue
"", // consumer tag (empty = auto-generated)
true, // auto-ack
false, // exclusive
false, // no-local
false, // no-wait
nil,
)
if err != nil {
log.Fatalf("Failed to register consumer: %v", err)
}
go func() {
for d := range msgs {
log.Printf("[Fanout-1] Received: %s", d.Body)
}
}()
}
同理,HandleMessageFanout2 应使用 "queue-fanout-2" 作为队列名,并完成相同声明与绑定流程。
✅ 补充:生产者示例(Java 或 Go)需向 logs 交换器发布
// 示例:Go 生产者(需在同 channel 上)
err := ch.Publish(
"logs", // exchange
"", // routing key (ignored)
false, // mandatory
false, // immediate
amqp.Publishing{
ContentType: "text/plain",
Body: []byte("Hello from Fanout!"),
})
? 总结与注意事项
- ? Fanout Exchange 不依赖 routing key,所有绑定队列无条件接收全部消息;
- ? 每个消费者必须对应独立队列(不能共用 queue name),否则退化为竞争消费;
- ?️ 建议将 ExchangeDeclare 和 QueueDeclare 放在应用启动时集中初始化,避免重复声明;
- ? 若测试中旧队列残留影响行为,可通过 RabbitMQ Management UI 清理或使用 autoDelete: true(仅用于开发);
- ? 官方权威参考:RabbitMQ Tutorial 3 — Publish/Subscribe (Go)。
遵循以上结构,两个 Go 消费者将同时、独立、完整地收到每一条 fanout 消息,真正实现广播语义。










