RabbitMQ Priority Queue not working in Golang

Viewed 135

I am using RabbitMQ with Golang. I've declared the queue in following way:

func declareTargetQueue(ch *amqp.Channel) (amqp.Queue, error) {
   args := make(amqp.Table)
   args["x-max-priority"] = uint8(9)
   return ch.QueueDeclare(
       "my-queue", // name
       true,                          // durable
       false,                         // delete when unused
       false,                         // exclusive
       false,                         // no-wait
       args,                          // arguments
   )
}

Here is my Producer function:

func (cli *DotRabbitCli) Produce(ctx *context.Context, request DotRabbitProduceRequest) {
//set connection
//create channel and assign it the name of 'ch'
//create headers object of type amqp.Table

err = ch.Publish(
    request.Exchange, // exchange
    request.Queue,    // routing key
    false,            // mandatory
    false,            // immediate
    amqp.Publishing{
        DeliveryMode: amqp.Persistent,
        ContentType:  "application/json",
        Body:         request.Data,
        Headers:      headers,
        Priority:     request.MessagePriority,
    })

if err != nil {
    ...
}

}

Here's the consumer code:

err = ch.Qos(
        1,     // prefetch count
        0,     // prefetch size
        false, // global
    )
    if err != nil {
        ..
    }

    msgs, err := ch.Consume(
        targetQueue.Name, // queue
        "",               // consumer
        false,            // auto-ack
        false,            // exclusive
        false,            // no-local
        false,            // no-wait
        nil,              // args
    )
    if err != nil {
        ...
    }

    forever := make(chan bool)

    go func() {
        for d := range msgs {
            request := MyRequest{}
            e := json.Unmarshal(d.Body, &request)
            if e != nil {
                // log error
                d.Ack(false)
                continue
            }

            d.Ack(false)

            if err := doMyWork(); err != nil {
                ConsumerFailLog(ctx, request, err, "Worker_Failed")
            } else {
                ConsumerSuccessLog(ctx, request, "Worker_Success")
            }
        }
    }()

    log.Printf(" [*] Waiting for messages. To exit press CTRL+C")
    <-forever

Now, here's what happened:

  • For first 10 minutes, there were 50 messages published with priority 4.
  • After 10 minutes, while above 50 messages were in queue, I published another message with priority 9.
  • I waited for so long but none of the 51 messages got acknowledged by the consumer.
  • After sometime, the consumer was dead and all the messages remained in the queue.

Although I read the documentation here, but I'm not able to find a solution. Can anyone please tell me what I can do to resolve this scenario?

0 Answers
Related