发放优惠卷时增加了通知
This commit is contained in:
@@ -39,20 +39,19 @@ func NewRabbitMQClient() (*RabbitMQ, error) {
|
||||
}
|
||||
|
||||
// Publish 生产
|
||||
func (r *RabbitMQ) Publish(queueName, exchangeName string, message interface{}) error {
|
||||
q, err := r.channel.QueueDeclare(
|
||||
func (r *RabbitMQ) Publish(queueName, exchangeName, routingKey string, message interface{}) error {
|
||||
_, err := r.channel.QueueDeclare(
|
||||
queueName, // 队列名字
|
||||
false, // 消息是否持久化
|
||||
true, // 消息是否持久化
|
||||
false, // 不使用的时候删除队列
|
||||
false, // 排他
|
||||
false, // 是否等待服务器确认
|
||||
false, // 是否排他
|
||||
false, // 是否阻塞
|
||||
nil, // arguments
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Serialize the message
|
||||
body, err := json.Marshal(message)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -60,7 +59,7 @@ func (r *RabbitMQ) Publish(queueName, exchangeName string, message interface{})
|
||||
|
||||
err = r.channel.Publish(
|
||||
exchangeName, // exchange(交换机名字)
|
||||
q.Name, //
|
||||
routingKey, //
|
||||
false, // 是否为无法路由的消息进行返回处理
|
||||
false, // 是否对路由到无消费者队列的消息进行返回处理 RabbitMQ 3.0 废弃
|
||||
amqp.Publishing{
|
||||
@@ -74,9 +73,9 @@ func (r *RabbitMQ) Publish(queueName, exchangeName string, message interface{})
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *RabbitMQ) PublishWithDelay(exchangeName, routingKey string, message interface{}, delay time.Duration) error {
|
||||
func (r *RabbitMQ) PublishWithDelay(queueName, exchangeName, routingKey string, message interface{}, delay time.Duration) error {
|
||||
err := r.channel.ExchangeDeclare(
|
||||
exchangeName, // name
|
||||
queueName, // name
|
||||
"x-delayed-message", // type
|
||||
true, // durable
|
||||
false, // auto-deleted
|
||||
|
||||
Reference in New Issue
Block a user