Celery, commit-aware enqueue, retries, идемпотентность и Django Tasks
Статус: примеры рассчитаны на Django 5.2/6.0 и актуальную ветку Celery 5.6. В Celery 5.4 появился
delay_on_commit(). Django 6.0 добавил собственный Tasks API, но не production worker.
Фоновая задача переносит работу за пределы HTTP request-response. Пользователь получает короткий ответ, а отдельный worker позже строит отчёт, обрабатывает файл или обращается к интеграции.
Это не просто способ «ускорить всё дольше 200 мс». Перенос меняет контракт:
Поэтому фоновой делают работу, которую действительно можно завершить позже и для которой спроектированы повтор, идемпотентность и наблюдаемость.
Django producer -> broker -> Celery worker -> внешняя система/БД
^
|
Celery BeatBroker и result backend — разные роли, даже если оба используют Redis. RabbitMQ и Redis имеют разные модели доставки, persistence и эксплуатации; выбор нельзя свести к «один всегда надёжнее». Зафиксируйте требования к потере данных, redelivery, latency, размеру сообщений и восстановлению после сбоя, затем проверьте настройки конкретного broker.
Result backend не нужен задаче только ради того, чтобы она выполнилась. Если результат не читается, включайте ignore_result=True и наблюдайте outcome через бизнес-запись, логи и метрики.
# config/celery.py
import os
from celery import Celery
os.environ.setdefault("DJANGO_SETTINGS_MODULE", "config.settings")
app = Celery("config")
app.config_from_object("django.conf:settings", namespace="CELERY")
app.autodiscover_tasks()# config/__init__.py
from .celery import app as celery_app
__all__ = ("celery_app",)# settings.py
import os
CELERY_BROKER_URL = os.environ["CELERY_BROKER_URL"]
CELERY_ACCEPT_CONTENT = ["json"]
CELERY_TASK_SERIALIZER = "json"
CELERY_TIMEZONE = "Europe/Moscow"JSON уменьшает риск исполнения произвольного объекта по сравнению с недоверенным pickle payload и заставляет сохранять явный контракт сообщения.
Worker и scheduler запускаются отдельными управляемыми процессами:
celery -A config worker --loglevel=INFO
celery -A config beat --loglevel=INFOКомбинация worker -B удобна лишь для локальной разработки. В production Beat запускают отдельно и обычно в одном экземпляре для одного расписания, иначе одна периодическая запись может публиковаться несколько раз.
from celery import shared_task
@shared_task(ignore_result=True)
def rebuild_search_document(article_id):
article = Article.objects.get(pk=article_id)
search_index.upsert(article.to_search_document())Передавайте JSON-примитивы: ID, строки, числа, списки и словари с явной версией. ORM instance не является стабильным сообщением и обычно не сериализуется JSON.
У ID есть важная семантика: задача прочитает состояние на момент выполнения, а не публикации. Для письма-счёта или юридического документа это может быть неверно. Тогда создайте immutable snapshot/outbox row в транзакции и передайте его ID. Не передавайте огромный файл через broker — положите его в storage и передайте ссылку/идентификатор с permission и TTL.
Версионируйте долгоживущие сообщения, если во время rolling deployment старые workers и новый producer могут работать одновременно. Worker должен уметь прочесть прежний формат до опустошения очереди.
with transaction.atomic():
report = Report.objects.create(owner=request.user)
build_report.delay(report.pk)Worker может стартовать до commit и не найти Report. Хуже того, транзакция может откатиться, а задача уже существует.
Универсальное решение Django:
from functools import partial
from django.db import transaction
with transaction.atomic():
report = Report.objects.create(owner=request.user)
transaction.on_commit(partial(build_report.delay, report.pk))Начиная с Celery 5.4, обычная Django-интеграция предоставляет shortcut:
build_report.delay_on_commit(report.pk)delay_on_commit() не возвращает task ID: сообщение появится только после commit. Если task ID нужен HTTP-ответу, создайте собственный ReportJob с UUID внутри транзакции и возвращайте ID этой бизнес-записи.
При custom base class нужно наследоваться от celery.contrib.django.task.DjangoTask, иначе shortcut может отсутствовать.
on_commit() закрывает гонку с невидимыми строками, но между commit базы и успешной отправкой в broker остаётся окно сбоя. Если потеря задачи недопустима, используйте transactional outbox:
Это добавляет код, зато устраняет попытку сделать атомарной операцию сразу в двух независимых системах.
Практический дизайн Celery должен допускать повторное выполнение. Worker может завершить внешнюю операцию и упасть до acknowledgment. Broker доставит сообщение снова. Сетевой timeout тоже не говорит, выполнил ли поставщик запрос.
Настройки acks_late, task_reject_on_worker_lost, visibility timeout и prefetch влияют на момент redelivery и пропускную способность, но не создают exactly-once side effect. Меняйте их только вместе с идемпотентной задачей и тестом отказа.
Надёжные приёмы:
UPDATE ... WHERE status = ... для перехода состояния;update_or_create() там, где его lookup действительно уникален;Один task_id в Redis и временный cache lock не доказывают идемпотентность: ключ может исчезнуть, а task ID измениться при повторной публикации.
Для email есть неустранимое окно: SMTP/provider принял письмо, worker упал до записи sent_at, затем retry отправил его снова. Если provider поддерживает idempotency key — используйте его. Иначе выбирайте осознанный компромисс и отслеживайте дубли.
from celery import shared_task
@shared_task(
autoretry_for=(CrmTimeout, CrmUnavailable),
retry_backoff=5,
retry_backoff_max=300,
retry_jitter=True,
max_retries=6,
ignore_result=True,
)
def sync_order(order_id):
order = Order.objects.get(pk=order_id)
crm_client.upsert_order(
order.to_crm_payload(),
idempotency_key=f"order:{order.pk}:v{order.revision}",
timeout=10,
)Повторяют timeout, временную недоступность и rate-limit с корректной задержкой. Не повторяют бесконечно ошибку валидации, отсутствие обязательного объекта или программный AttributeError. autoretry_for=(Exception,) превращает постоянную ошибку кода в очередь одинаковых сбоев.
Timeout внешнего клиента обязателен независимо от Celery time limit. Soft/hard time limits — последняя страховка, их поддержка и поведение зависят от worker pool и платформы. Hard kill не выполняет finally надёжно, поэтому cleanup и согласованность нельзя строить только на нём.
Celery state PENDING может означать «неизвестный task ID», «ещё не началась» или «result backend не хранит запись». Пользовательскому интерфейсу лучше показывать собственную модель:
ReportJob: queued -> running -> succeeded
\-> failedПереходы делайте атомарными, храните число попыток, безопасное сообщение об ошибке и timestamps. Task может обновлять progress, но слишком частые записи сами становятся нагрузкой.
Не возвращайте traceback клиенту. Для повторного запуска создавайте новую попытку или явно переводите допустимый terminal state, сохраняя аудит.
from celery.schedules import crontab
CELERY_BEAT_SCHEDULE = {
"expire-invitations": {
"task": "accounts.tasks.expire_invitations",
"schedule": crontab(minute="*/15"),
},
}Beat публикует сообщения по расписанию, но не гарантирует, что предыдущий запуск уже завершился. Периодическая задача должна терпеть overlap либо сама получать распределённый lease/DB-lock. Для расписания из admin используют django-celery-beat, но его database scheduler и миграции становятся частью эксплуатации.
Timezone, переходы DST и catch-up после простоя требуют отдельного решения. Для ежедневного бизнес-отчёта храните логическую дату запуска и unique constraint, а не полагайтесь только на wall clock.
Разделяйте workload, когда один тип мешает другому:
CELERY_TASK_ROUTES = {
"mail.tasks.*": {"queue": "mail"},
"exports.tasks.*": {"queue": "exports"},
}celery -A config worker -Q mail --concurrency=8
celery -A config worker -Q exports --concurrency=2Короткие уведомления не должны стоять за часовым экспортом. CPU-bound workers и I/O-bound workers получают разные concurrency, memory limit, autoscaling и time limits. Priority конкретного broker может иметь ограничения; отдельная очередь часто предсказуемее.
Минимальные сигналы:
Flower удобен для оперативного обзора, celery inspect — для диагностики, error tracker — для исключений. Ни один из них не заменяет алерты и бизнес-метрику «отчёт пользователя не готов 20 минут».
Не помещайте персональные данные и секреты в task name, args representation и логи. Event monitoring Celery способен показывать аргументы.
На уровне unit вызывайте чистую доменную функцию без broker. Отдельно тестируйте task wrapper, retry classification и идемпотентность.
Eager mode удобен, но скрывает JSON serialization, commit race, redelivery и отличия worker process. Критичные потоки проверяйте integration-тестом с тем же типом broker и настоящим worker:
Django 6.0 добавил django.tasks: декоратор @task, enqueue()/aenqueue(), JSON-валидацию аргументов, task backends и результаты.
from django.tasks import task
@task(queue_name="mail")
def send_digest(user_id):
# Прикладная работа.
...
send_digest.enqueue(user_id=42)Встроенные ImmediateBackend и DummyBackend предназначены для разработки и тестов. Django не поставляет production worker: нужен сторонний backend/исполнитель. Celery остаётся зрелой системой с workers, Beat, routing, canvas и большой экосистемой.
Django Tasks не отменяет транзакционную границу: enqueue() тоже регистрируют через transaction.on_commit(). Выбирайте abstraction после оценки production backend, retries, schedule, мониторинга и миграции существующих задач, а не только по удобству декоратора.
delay() сразу после создания строки внутри atomic() заменён на delay_on_commit()/on_commit().task_id не считается гарантией идемпотентности.autoretry_for=(Exception,) не рекомендуется как универсальный пример: retry нужен временным, классифицированным сбоям.Далее: Docker и Деплой