登录代码提交
This commit is contained in:
54
trunk/center/common/rebbitMQ/config.go
Normal file
54
trunk/center/common/rebbitMQ/config.go
Normal file
@@ -0,0 +1,54 @@
|
||||
package rebbitMQ
|
||||
|
||||
import (
|
||||
configYaml "common/configsYaml"
|
||||
"github.com/streadway/amqp"
|
||||
"goutil/logUtil"
|
||||
)
|
||||
|
||||
var RabbitMQConn *amqp.Connection
|
||||
var RabbitMQChannel *amqp.Channel
|
||||
|
||||
// 初始化rabbitMQ
|
||||
func init() {
|
||||
|
||||
rabbitMQAddress := configYaml.GetRabbitMQAddress()
|
||||
|
||||
//是否有mq配置
|
||||
if rabbitMQAddress == "" {
|
||||
return
|
||||
}
|
||||
|
||||
// 连接到 RabbitMQ 服务器
|
||||
var err error
|
||||
RabbitMQConn, err = amqp.Dial(rabbitMQAddress)
|
||||
if err != nil {
|
||||
|
||||
//抛出一个异常
|
||||
logUtil.FatalLog("Failed to connect to RabbitMQ: %s,err:%s", rabbitMQAddress, err.Error())
|
||||
}
|
||||
|
||||
// 打开一个通道
|
||||
RabbitMQChannel, err = RabbitMQConn.Channel()
|
||||
if err != nil {
|
||||
|
||||
//抛出一个异常
|
||||
logUtil.FatalLog("Failed to open a channel,err:%s", err.Error())
|
||||
}
|
||||
|
||||
//队列名称
|
||||
queueName := configYaml.GetRabbitMQName()
|
||||
|
||||
// 声明一个队列
|
||||
_, err = RabbitMQChannel.QueueDeclare(
|
||||
queueName, // 队列名称
|
||||
true, // 是否持久化
|
||||
false, // 是否在使用后删除
|
||||
false, // 是否排他
|
||||
false, // 是否阻塞
|
||||
nil, // 其他参数
|
||||
)
|
||||
if err != nil {
|
||||
logUtil.FatalLog("Failed to declare a queue,queueName:%s,err:%s", queueName, err.Error())
|
||||
}
|
||||
}
|
||||
25
trunk/center/common/rebbitMQ/sendData.go
Normal file
25
trunk/center/common/rebbitMQ/sendData.go
Normal file
@@ -0,0 +1,25 @@
|
||||
package rebbitMQ
|
||||
|
||||
import (
|
||||
configYaml "common/configsYaml"
|
||||
"github.com/streadway/amqp"
|
||||
"goutil/logUtilPlus"
|
||||
)
|
||||
|
||||
// 发送单挑mq消息
|
||||
func sendMqData(data string) {
|
||||
|
||||
// 发布消息到队列
|
||||
err := RabbitMQChannel.Publish(
|
||||
"", // 交换机名称
|
||||
configYaml.GetRabbitMQName(), // 路由键
|
||||
false, // 是否强制
|
||||
false, // 是否立即
|
||||
amqp.Publishing{
|
||||
ContentType: "text/plain",
|
||||
Body: []byte(data),
|
||||
})
|
||||
if err != nil {
|
||||
logUtilPlus.ErrorLog("Failed to publish a message,err:%s", err.Error())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user