一、RabbitMQ 基础概念AMQP 协议AMQP高级消息队列协议与语言、平台无关的消息协议RabbitMQ 基于该协议。Spring AMQP 简介Spring AMQP 是基于 AMQP 协议封装的 API 规范分为两部分spring-amqp基础抽象层spring-rabbit底层 RabbitMQ 具体实现二、Go 客户端安装注意这里有两套 Go 客户端新版官方rabbitmq-amqp-go-clientAMQP 1.0社区常用streadway/amqpAMQP 0-9-1下面 Demo 用的这个学习推荐# 新版官方客户端安装 go get github.com/rabbitmq/rabbitmq-amqp-go-client导入import ( context log rmq github.com/rabbitmq/rabbitmq-amqp-go-client/pkg/rabbitmqamqp )下面示例代码使用的是streadway/amqpAMQP0-9-1目前学习最简单和 SpringAMQP 底层协议一致三、简单模型生产者直接发消息到队列Simple 模型流程Publisher → Queue → Consumer 案例simple.queue在 RabbitMQ 控制台创建队列simple.queue代码里QueueDeclare是幂等不存在自动创建生产者发送消息到队列消费者监听队列接收消息SpringAMQP 开发步骤对照方便 Java 转 Go 理解引入 RabbitMQ 依赖配置 RabbitMQ 服务信息host、port、virtual-host、username、password发送消息使用RabbitTemplate消费消息使用RabbitListener注解监听队列Go 生产者代码send.go等价 SpringrabbitTemplate.convertAndSend(queueName, msg)package main import ( fmt github.com/streadway/amqp ) func main() { // 1. 连接 rabbitmq conn, err : amqp.Dial(amqp://guest:guestlocalhost:5672/) if err ! nil { panic(err) } defer conn.Close() // 2. 创建通道 ch, err : conn.Channel() if err ! nil { panic(err) } defer ch.Close() queueName : simple.queue // 3. 声明队列幂等。队列存在则不操作不存在自动创建 _, err ch.QueueDeclare( queueName, false, // durable是否持久化 false, // autoDelete是否自动删除 false, // exclusive独占队列 false, // noWait nil, // args额外参数 ) if err ! nil { panic(err) } msg : hello, amqp! // 4. 发布消息 err ch.Publish( , // exchange空默认交换机和Spring行为一致 queueName, // routingKey直接填队列名消息直接路由到此队列 false, // mandatory false, // immediate amqp.Publishing{ ContentType: text/plain, Body: []byte(msg), }, ) if err ! nil { panic(err) } fmt.Printf(发送消息成功%s\n, msg) }Go 消费者代码consumer.go等价 SpringRabbitListener(queues simple.queue)持续监听队列package main import ( fmt github.com/streadway/amqp log ) func main() { // 1. 连接RabbitMQ conn, err : amqp.Dial(amqp://guest:guestlocalhost:5672/) if err ! nil { log.Fatalf(连接失败: %v, err) } defer conn.Close() // 2. 创建通道 ch, err : conn.Channel() if err ! nil { log.Fatalf(打开通道失败: %v, err) } defer ch.Close() queueName : simple.queue // 3. 声明队列幂等 _, err ch.QueueDeclare( queueName, false, false, false, false, nil, ) if err ! nil { log.Fatalf(声明队列失败: %v, err) } // 4. 注册消费者开始消费消息 msgs, err : ch.Consume( queueName, // 监听队列名 , // consumer true, // autoAcktrue收到消息自动确认和Spring默认近似 false, // exclusive false, // noLocal false, // noWait nil, // args ) if err ! nil { log.Fatalf(注册消费失败: %v, err) } fmt.Println(Go消费者启动监听 simple.queue ...) // 循环读取消息持续处理消息 for d : range msgs { msg : string(d.Body) fmt.Printf(go 消费者接收到消息[%s]\n, msg) fmt.Println(消息处理完成) } }四、核心要点总结连接地址不要写错amqp://账号:密码ip:5672/QueueDeclare幂等特性重复调用不会报错不存在则新建队列Simple 模型空交换机routingKey 直接写队列名消息直接投递队列autoAck自动确认消息收到后 MQ 立刻删除消息生产环境一般关闭手动 ACK 保证可靠性SpringAMQP 和 Go streadway/amqp 底层都是 AMQP0-9-1可以对照理解学习
