Found strange behavior implementing tasks pipeline on Celery. Most of the time tasks chain executed, but sometimes it just silently stops in the middle after successful run of the previous task.
Example pipeline
pipelines = []
for task in Factory.gen_tasks(request):
pipelines.append(
chain(
first_process_step.si(task),
second_process_step.si(task),
group(
chain(
post_process_first_step.s(task.id),
post_process_second_step.s(task.id),
),
notify_user.s(task.id),
),
)
)
async_task = group(pipelines).apply_async()
Some config options:
celery_app.conf.update(task_acks_late=True)
celery_app.conf.update(task_reject_on_worker_lost=True)
celery_app.conf.update(worker_proc_alive_timeout=20)
celery_app.conf.update(worker_lost_wait=10)
So most of the time (99%, probably) everything is fine and task executed after another task. But sometimes execution stopped after first_process_step or sometimes after second_process_step
In both cases, I see in logs that the previous task completed, but the next step task is not received by a worker. Also, the queue is empty. So I started to think that a message was not sent to the broker.
App deployed to Kube and works with AWS Redis and RabbitMQ.
UPDATE:
So, if anyone will have similar behavior when the pipeline looks pretty regular, you use RabbitMQ and sometimes messages lost silently then you need to consider adding this not documented option:
celery_app.conf.update(
broker_transport_options={"confirm_publish": True},
)
Which enabled ConfirmPublish and ensures that a message will be processed at least once.
Defect about documentation update: https://github.com/celery/celery/issues/5410