前言: 刚开始学习 RabbitMQ 时,队列、交换机都是在 MQ 的 Web 控制台手动创建。但在实际开发中,业务队列数量很多,不可能每次都手动在 RabbitMQ 控制台创建交换机和队列。推荐在代码中完成队列、交换机、绑定关系的声明,程序启动自动校验 / 创建。
一、SpringAMQP(Java)声明方式
SpringAMQP 提供了三个核心类,用来声明队列、交换机以及二者的绑定关系,并且配套链式工厂构建器:
Queue:用于声明队列,可以使用工厂类QueueBuilder构建Exchange:用于声明交换机,可以使用工厂类ExchangeBuilder构建Binding:用于声明队列和交换机的绑定关系,可以使用工厂类BindingBuilder构建
// 链式工厂构建示例 Queue queue = QueueBuilder.durable("fanout.queue1").build(); Exchange exchange = ExchangeBuilder.fanoutExchange("hmall.fanout").build(); Binding binding = BindingBuilder.bind(queue).to(exchange);底层原理:Spring 将 AMQP 协议指令封装为 Bean 对象,容器启动时自动执行协议指令,完成交换机、队列、绑定的创建。我们先在内存构造对象,由 Spring 统一执行。
二、Go streadway/amqp 声明方式
Go 的streadway/amqp原生库没有提供 QueueBuilder / ExchangeBuilder / BindingBuilder 工厂类,也不存在 Binding 结构体来保存绑定关系。 Go 库是直接操作 AMQP 底层协议:调用方法的一瞬间,就向 RabbitMQ 服务端发送协议指令,直接执行创建 / 校验操作。
核心三个方法:
ch.QueueDeclare():执行协议指令,在 RabbitMQ 服务端创建 / 校验队列ch.ExchangeDeclare():执行协议指令,在 RabbitMQ 服务端创建 / 校验交换机ch.QueueBind():执行协议指令,在 RabbitMQ 服务端建立队列与交换机的绑定关系
Go 完整代码示例(Fanout 交换机 + 双队列 + 绑定)
package main import ( "log" "github.com/streadway/amqp" ) func failOnError(err error, msg string) { if err != nil { log.Fatalf("%s: %s", msg, err) } } func main() { // 连接RabbitMQ conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/") failOnError(err, "连接MQ失败") defer conn.Close() ch, err := conn.Channel() failOnError(err, "打开通道失败") defer ch.Close() // ========== 1. 声明 Fanout 交换机 hmall.fanout ========== err = ch.ExchangeDeclare( "hmall.fanout", // 交换机名称 amqp.ExchangeFanout, // 交换机类型 fanout false, // durable 持久化 false, // autoDelete false, // internal false, // noWait nil, // arguments ) failOnError(err, "声明 fanout 交换机失败") // ========== 2. 声明队列 fanout.queue1 ========== q1, err := ch.QueueDeclare( "fanout.queue1", //队列名 false, false, false, false, nil, ) failOnError(err, "声明队列 fanout.queue1 失败") // ========== 3. 绑定队列1 到交换机 ========== err = ch.QueueBind( q1.Name, "", // fanout模式 bindingKey为空,不生效 "hmall.fanout", false, nil, ) failOnError(err, "绑定 fanout.queue1 失败") // ========== 4. 声明队列 fanout.queue2 ========== q2, err := ch.QueueDeclare( "fanout.queue2", false, false, false, false, nil, ) failOnError(err, "声明队列 fanout.queue2 失败") // ========== 5. 绑定队列2 到交换机 ========== err = ch.QueueBind( q2.Name, "", "hmall.fanout", false, nil, ) failOnError(err, "绑定 fanout.queue2 失败") log.Println("交换机、队列、绑定全部创建完成!") }三、Java SpringAMQP vs Go streadway/amqp 核心对比
表格
| 对比项 | SpringAMQP(Java) | streadway/amqp(Go) |
|---|---|---|
| Builder 工厂 | 提供QueueBuilder/ExchangeBuilder/BindingBuilder链式构建器 | 无原生 Builder,需要自己封装 |
| 编程思想 | 先在内存构造 Queue/Exchange/Binding 对象,Spring 容器启动后统一执行声明 | 直接调用QueueDeclare/ExchangeDeclare/QueueBind,调用即发送 AMQP 指令到 MQ 服务端 |
| Binding 对象 | Binding 是独立对象,描述队列交换机绑定关系 | 没有 Binding 结构体,绑定动作通过QueueBind方法直接执行 |
| 抽象层级 | 高层封装,面向对象 | 底层 AMQP 协议 SDK,偏原生 |
一句话总结: SpringAMQP 是先构造对象,交给容器执行;Go streadway/amqp 是调用方法就直接操作 MQ 服务端。
四、补充知识点
ExchangeDeclare、QueueDeclare是幂等操作: 如果交换机 / 队列已经存在,不会重复新建,只会校验参数是否匹配;参数不一致会直接报错,这也是代码声明方式的优势。