Recently we encountered an issue with Celery where, when one worker X has network connectivity issues the tasks appears to be redelivered to other workers. There's nothing wrong with this, but the problem is that very few workers are suddenly prefetching like 150 of these redelivered tasks. This happens even when we set worker_prefetch_multiplier to 1
So when we re-launch the worker X which had connectivity issues, it is not fetching any tasks anymore - because they got prefetched by other workers. Even though we didn't want the other workers to prefetch more than 1 task at the time.
The prefetching works fine on regular task submission and processing, however in situation like described above - where one of the workers has connectivity issues - weird things start to happen.
We know that the other workers are absorbing most of these messages because we tried to stop all of them. Upon pressing CTRL+C we see message on two workers like these:
Restoring 150 unacknowledged message(s)
instead of usual message (when things are fine)
Restoring 1 unacknowledged message(s)
How is this possible that it prefetched 150 tasks upon other worker failure, when we set worker_prefetch_multiplier to 1?
We inspected the redis once the tasks has been restored to redis queue to see what we can find there. Not a celery expert here but this celery internals line seemed interesting
\", \"kwargsrepr\": \"{}\", \"origin\": \"WORKER_NAME_WHICH_HAD_CONNECTIVITY_ISSUES\", \"ignore_result\": false, \"redelivered\": true},
Origin pointed to the worker which had connection issues and redelivered is set to true.
Does anyone have idea why suddenly few of the workers had 150 unacknowledged tasks each, even when we had worker_prefetch_multiplier set to 1?
This happened only on other worker network issues. In normal conditions the prefetching appears to work fine.
Below is shown how we launch worker and tasks:
# How the worker is started
celery -A tasks worker -n "someworkername" --loglevel=INFO -Q "hardcodedqueuename" -c 1
app.conf['worker_prefetch_multiplier'] = 1
app.conf.broker_transport_options = {"visibility_timeout": 6 * 3600}
signatures=[]
signatures.append(scan_task.s(arg1, args).set(queue=queue_name))
finalize_func = finalize.s(context)
group_task = group(signatures,
queue=queue_name) | finalize_func
group_task()