package connectors
import (
"../commons"
"github.com/streadway/amqp"
log "github.com/sirupsen/logrus"
"bytes"
"strconv"
"fmt"
)
func failOnError(err error, msg string) {
if err != nil {
log.Fatalf("%s: %s", msg, err)
panic(fmt.Sprintf("%s: %s", msg, err))
}
}
func rabbitConnector() (*amqp.Connection, error) {
rabbitConfig := commons.RabbitConfig
addr := bytes.Buffer{}
addr.WriteString("amqp://")
addr.WriteString(rabbitConfig.Username)
addr.WriteString(":")
addr.WriteString(rabbitConfig.Password)
addr.WriteString("@")
addr.WriteString(rabbitConfig.Host)
addr.WriteString(":")
addr.WriteString(strconv.Itoa(rabbitConfig.Port))
addr.WriteString("/")
addr.WriteString(rabbitConfig.Vhost)
conn, err := amqp.Dial(addr.String())
if err != nil {
failOnError(err, "Failed to connect to RabbitMQ")
}
return conn, err
}
func Send(msg string) {
defer func() {
if err := recover(); err != nil {
log.Errorf("RabbitMQ发送存储消息错误 %s", err)
}
}()
rabbitConfig := commons.RabbitConfig
conn, err := rabbitConnector()
failOnError(err, "Failed to connect to RabbitMQ")
defer conn.Close()
ch, err := conn.Channel()
failOnError(err, "Failed to open a channel")
defer ch.Close()
err = ch.ExchangeDeclare(
rabbitConfig.Exchange, // name
"topic", // type
true, // durable
false, // auto-deleted
false, // internal
false, // no-wait
nil, // arguments
)
failOnError(err, "Failed to declare an exchange")
err = ch.Publish(
rabbitConfig.Exchange, // exchange
"custom", // routing key
false, // mandatory
false, // immediate
amqp.Publishing{
ContentType: "text/plain",
Body: []byte(msg),
})
failOnError(err, "Failed to publish a message")
}
func RabbitConsume() {
conn, err := amqp.Dial(fmt.Sprintf("amqp://%s:%s@%s:%d/dashboard", commons.RabbitConfig.Username, commons.RabbitConfig.Password, commons.RabbitConfig.Host, commons.RabbitConfig.Port))
failOnError(err, "Failed to connect to RabbitMQ")
defer conn.Close()
ch, err := conn.Channel()
failOnError(err, "Failed to open a channel")
defer ch.Close()
err = ch.ExchangeDeclare(
"dashboard", // name
"topic", // type
true, // durable
false, // auto-deleted
false, // internal
false, // no-wait
nil, // arguments
)
failOnError(err, "Failed to declare an exchange")
q, err := ch.QueueDeclare(
"", // name
true, // durable
false, // delete when usused
false, // exclusive
false, // no-wait
nil, // arguments
)
failOnError(err, "Failed to declare a queue")
err = ch.QueueBind(
q.Name, // queue name
"custom", // routing key
"dashboard", // exchange
false,
nil)
failOnError(err, "Failed to bind a queue")
msgs, err := ch.Consume(
q.Name, // queue
"", // consumer
true, // auto-ack
false, // exclusive
false, // no-local
false, // no-wait
nil, // args
)
failOnError(err, "Failed to register a consumer")
forever := make(chan bool)
go func() {
for d := range msgs {
log.Printf(" [x] %s", d.Body)
}
}()
log.Printf(" [*] Waiting for logs. To exit press CTRL+C")
<-forever
}
RabbitMQ消费 发送
©著作权归作者所有,转载或内容合作请联系作者
- 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
- 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
- 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
推荐阅读更多精彩内容
- 文档: https://www.rabbitmq.com/confirms.html 介绍:使用像RabbitMQ...
- 在上一章第四十一章: 基于SpringBoot & RabbitMQ完成DirectExchange分布式消息消费...
- 分布式系统中,我们广泛运用消息中间件进行系统间的数据交换,便于异步解耦。现在开源的消息中间件有很多,前段时间我们自...
- http://www.jianshu.com/p/4112d78a8753 接这篇 在上文中,主要实现了可靠模式的...