How to unacknowledge message when I run a celery task?

Viewed 946

Now I have some sync jobs, which are stateful, so if the task failed, I must unacknowledge message then let them go to the front of RabbitMQ. But when I try to raise an error, I found celery still acknowledges this message, and the queue has been cleared.

@celery.task(bind=True)
def my_task(self, *args, **kwargs):
    raise ValueError

And I found celery task has a method called retry, but it will add the task to the back of the queue. This is not what I want.

@celery.task(bind=True)
def my_task(self, *args, **kwargs):
    try:
        raise ValueError
    except Exception:
        self.retry(countdown=15)

Even I can't do that with a kill signal:

os.kill(os.getpid(), signal.SIGKILL)

What should I do? Did celery provide some error so I can raise this error to notify celery don't acknowledge my message?

2 Answers

As per documentation celery is capable of RabbitMQs priority queues. https://docs.celeryproject.org/en/latest/faq.html#does-celery-support-task-priorities

Therefore you should be able to push the retried tasks to the front of your queue by prioritizing them higher than your regular tasks.

@celery.task(bind=True)
def my_task(self, *args, **kwargs):
    try:
        raise ValueError
    except Exception:
        self.retry(countdown=15, priority=9)

or

According to this github issue, you can also assign the retried task to a new dedicated queue and assign your ressources to prioritise this queue.

@celery.task(bind=True)
def my_task(self, *args, **kwargs):
    try:
        raise ValueError
    except Exception:
        self.retry(countdown=15, queue='prioritized_queue_name') 
Related