Почти в каждом приложении есть задача: накапливать записи в таблице и обрабатывать их в фоне. Платежи, уведомления, отчёты — что угодно. Первое решение пишется быстро: расписание, выборка пачки, цикл. Оно даже работает — до первого деплоя с двумя репликами или до первой ошибки во внешней системе.
Эта статья объясняет, что именно ломается и как сделать фоновую обработку надёжно.
Типичный первый подход
Вот код, который встречается в большинстве проектов:
@app.task
def run_payment_batch() -> None:
with session_factory() as session:
batch = session.execute(
select(Payment)
.where(Payment.status == PaymentStatus.NEW)
.order_by(Payment.created_at)
.limit(100)
).scalars().all()
for payment in batch:
try:
payment_service.process(payment.id)
except Exception:
logger.warning("Failed: %s", payment.id, exc_info=True)
app.conf.beat_schedule = {
"payment-batch": {"task": "run_payment_batch", "schedule": 300.0},
}
Выглядит разумно: раз в 5 минут берём 100 записей со статусом NEW и обрабатываем. Но у этого кода есть несколько серьёзных проблем.
Что здесь ломается
Очередь застревает из-за одной битой записи
Выборка сортируется по дате создания: ORDER BY created_at ASC. Это значит, что при следующем запуске обработчик возьмёт те же самые записи, что и в прошлый раз.
Если один платёж падает с ошибкой каждый раз — он никуда не уходит. Он остаётся в начале очереди и вместе с сотней таких же «застрявших» записей блокирует все новые. Новые платежи накапливаются, но до обработки не доходят.
Счётчика попыток нет. Отложенного повтора нет. Терминального статуса «не вышло» нет.
Зависшие записи после перезапуска
Обработчик берёт запись и сразу ставит ей статус PROCESSING. Потом идёт во внешнюю систему. Если в этот момент сервис упал — обновление статуса в обратно в NEW уже не происходит. Запись остаётся в PROCESSING навсегда.
Обработчик выбирает только NEW, поэтому такую «зависшую» запись он больше не возьмёт. После каждого перезапуска под нагрузкой таких записей накапливается несколько — и никто о них не знает до жалобы от клиента.
Двойная обработка при двух репликах
Запустили два экземпляра сервиса — оба обработчика читают одни и те же 100 строк одновременно. Один платёж отправится дважды. Для уведомлений это неудобство, для платежей — серьёзный инцидент.
Никакой блокировки выборки нет, и несколько параллельных воркеров договориться между собой не могут.
Один зависший вызов останавливает всё
Обработка идёт последовательно: один платёж → запрос во внешнюю систему → следующий платёж. Если внешняя система «задумалась» на одном запросе — весь цикл встал и ждёт. Таймаутов нет. Остальные 99 платежей стоят в очереди.
Как сделать правильно: очередь задач на PostgreSQL
Для большинства случаев отдельный брокер сообщений не нужен — база данных уже есть, транзакции уже работают. Нужна правильная механика поверх.
Запись захватывают атомарно, при успехе закрывают, при ошибке возвращают в очередь с задержкой, после предела попыток уводят в терминальный статус.
Шаг 1: добавить нужные поля
ALTER TABLE payment
ADD COLUMN attempts int NOT NULL DEFAULT 0,
ADD COLUMN next_attempt_at timestamptz NOT NULL DEFAULT now(),
ADD COLUMN claimed_until timestamptz;
attempts— сколько раз пытались обработать.next_attempt_at— когда можно попробовать снова (для повторных попыток с задержкой).claimed_until— до какого момента запись «занята» воркером.
Шаг 2: атомарно захватить пачку
Главная идея — захват и обработка должны быть разными транзакциями. Захват — короткая операция, которая атомарно помечает записи как «взяты в работу» прямо в запросе UPDATE:
CLAIM_BATCH_SQL = text("""
UPDATE payment SET
status = 'PROCESSING',
attempts = attempts + 1,
claimed_until = now() + interval '5 minutes'
WHERE id IN (
SELECT id FROM payment
WHERE (status = 'NEW' AND next_attempt_at <= now())
OR (status = 'PROCESSING' AND claimed_until < now())
ORDER BY created_at
LIMIT :limit
FOR UPDATE SKIP LOCKED
)
RETURNING id
""")
def claim_batch(limit: int) -> list[UUID]:
with session_factory.begin() as session:
rows = session.execute(CLAIM_BATCH_SQL, {"limit": limit})
return [row.id for row in rows]
Ключевая строка — FOR UPDATE SKIP LOCKED. Это PostgreSQL-конструкция: несколько воркеров одновременно запускают этот запрос, но каждый пропускает уже заблокированные строки и берёт следующие свободные. Двойной обработки не будет — база данных сама распределяет работу между воркерами без какого-либо координатора.
claimed_until решает проблему «зависших» записей: если воркер умер до завершения, через 5 минут условие claimed_until < now() включит эту запись обратно в выборку.
Шаг 3: обработка вне транзакции
Захватили записи, закоммитили — теперь идём во внешнюю систему. HTTP-вызов происходит между транзакциями, а не внутри:
@app.task
def run_payment_batch() -> None:
claimed = claim_batch(limit=100)
for payment_id in claimed:
process_one.delay(payment_id)
@app.task
def process_one(payment_id: UUID) -> None:
try:
decision = limits_client.check(payment_id)
finish_successfully(payment_id, decision)
except Exception as ex:
schedule_retry(payment_id, ex)
app.conf.beat_schedule = {
"payment-batch": {"task": "run_payment_batch", "schedule": 5.0},
}
Несколько важных изменений по сравнению с первым подходом:
- Записи обрабатываются параллельно — каждая уходит отдельной задачей в пул воркеров celery, а не последовательно в одном цикле. Один медленный вызов не блокирует остальные.
- Расписание теперь каждые 5 секунд, а не раз в 5 минут — пустой опрос стоит дёшево (один быстрый запрос), а задержка обработки падает с минут до секунд.
- HTTP-клиент
limits_clientдолжен иметь таймауты. Без явного таймаута один зависший вызов займёт воркера из пула навсегда.
Шаг 4: повторные попытки с задержкой и терминальный отказ
Когда обработка падает — запись не остаётся в начале очереди. Она уходит в конец с задержкой:
def schedule_retry(payment_id: UUID, ex: Exception) -> None:
with session_factory.begin() as session:
payment = session.get(Payment, payment_id)
if payment.attempts >= MAX_ATTEMPTS:
payment.mark_failed(str(ex)) # терминальный статус FAILED
else:
payment.schedule_next_attempt(backoff(payment.attempts)) # возврат в NEW с задержкой
backoff — функция, которая увеличивает паузу с каждой попыткой: например, 1 мин → 2 мин → 4 мин → 8 мин. После N попыток запись получает статус FAILED с причиной ошибки — это сигнал человеку разобраться вручную. Очередь не застревает.
Сам backoff считают с разбросом, иначе тысяча задач, упавших в одну секунду из-за лежащего соседа, вернётся в очередь в одну и ту же секунду и положит его снова:
import random
from datetime import timedelta
def backoff(attempts: int) -> timedelta:
base = timedelta(minutes=1 << min(attempts, 6)) # 1, 2, 4 … до 64 минут
jitter = timedelta(seconds=random.uniform(0, base.total_seconds() / 2))
return base + jitter
Число одновременных обработок задаёт пул воркеров, и он не бесконечен: восемь процессов означают не больше восьми открытых вызовов к соседу с этого узла, а task_acks_late вместе с worker_prefetch_multiplier=1 не даёт воркеру набрать задач впрок и потерять их при падении:
app.conf.update(
worker_concurrency=8, # восемь одновременных вызовов к соседу с одного узла
worker_prefetch_multiplier=1, # не брать задач впрок
task_acks_late=True, # подтверждать после выполнения, а не при получении
)
Идемпотентность внешнего вызова
Повторная попытка опасна, если первый вызов на самом деле дошёл: лимит проверен дважды, деньги списаны дважды. Поэтому ключ идемпотентности выводится из задачи и один и тот же на всех попытках, а не генерируется заново при каждом вызове:
@app.task
def process_one(payment_id: UUID) -> None:
task = load_task(payment_id)
try:
decision = limits_client.check(payment_id, idempotency_key=task.idempotency_key) # ключ хранится в задаче
finish_successfully(payment_id, decision)
except Exception as ex:
schedule_retry(payment_id, ex)
Сосед по этому ключу узнаёт повтор и возвращает прежний ответ вместо второго списания. Если у соседа нет поддержки ключа идемпотентности, повторять можно только запросы чтения, а запись после неизвестного исхода уходит человеку. Подробнее в статье про идемпотентность запросов в полёте.
Когда нужен фреймворк пакетной обработки
Airflow, Dagster, Prefect — это инструменты для конечных задач: выгрузить реестр платежей за месяц в файл, мигрировать десять миллионов строк, перетарифицировать всех клиентов раз в квартал.
Их ключевая возможность — перезапуск с места сбоя: «упало на 7-м миллионе — перезапустили, продолжило с 7-го». Для этого фреймворк хранит состояние прогона и поддерживает политики пропуска и повтора на уровне порций данных.
Для нашей задачи (непрерывная очередь входящих платежей) такой фреймворк избыточен: поток не имеет конца, перезапуск с места не нужен, а вопросы параллельного захвата и повторных попыток остаются теми же — фреймворк их не решает.
Простое правило: конечный пакет с перезапуском → фреймворк; бесконечный поток задач → очередь задач на PostgreSQL.
Когда добавлять брокер сообщений
Брокер (RabbitMQ, Kafka) решает те же задачи, что и очередь на PostgreSQL, но добавляет новые возможности: несколько независимых типов потребителей одного события, масштабирование воркеров без нагрузки на основную базу данных, встроенные очереди недоставленных сообщений.
Переходить к брокеру имеет смысл, когда объёмы вырастают до десятков тысяч задач в минуту или когда одно событие должны обрабатывать несколько разных сервисов. До этого — брокер лишняя инфраструктура с эксплуатационными расходами.
Глубже: ровно один раз в сутки: advisory lockрасширенное
FOR UPDATE SKIP LOCKED раздаёт очередь между репликами, и это решает «много задач, много воркеров». Задачу «сформировать реестр в три ночи ровно один раз» он не решает: очереди нет, есть расписание, и три реплики с планировщиком запустят её трижды.
Самый дешёвый замок уже есть в PostgreSQL: SELECT pg_try_advisory_lock(:key) возвращает true только одной сессии, остальные получают false и выходят. Ключ это число, например хеш имени задачи, замок живёт до pg_advisory_unlock или до конца сессии, так что упавший воркер его отпускает сам. Задание планировщика выглядит так: взять соединение, попробовать замок, при успехе выполнить и отпустить, при неудаче записать «пропущено, выполняет другая реплика» и завершиться. Соединение держат на всё время работы, иначе замок уйдёт вместе с ним в пул.
Второй вариант та же таблица задач: строка «реестр за 2026-10-03» с уникальным ключом по дате, которую вставляют ON CONFLICT DO NOTHING; кому вставка удалась, тот и выполняет. Он переживает рестарт в середине благодаря claimed_until и оставляет историю прогонов.
Третий вариант убирает проблему из приложения: CronJob в Kubernetes с concurrencyPolicy: Forbid запускает один под по расписанию. Для редких тяжёлых задач это честнее планировщика внутри сервиса, который делит процесс с обработкой запросов.
Глубже: наблюдаемость очереди задачрасширенное
Все ветки выше заканчиваются «сигнал человеку разобраться», и стоит назвать, по каким числам он приходит. У очереди на таблице их четыре, и все считаются одним запросом на расписании и отдаются как метрики.
Возраст самой старой необработанной задачи, now() - min(created_at) среди NEW: главный показатель, потому что глубина обманывает, а возраст нет. Тревога на возраст больше нормального интервала в несколько раз.
Число задач по статусам: NEW, PROCESSING, FAILED за последний час. Рост FAILED это сломанный сосед или сломанная логика, рост PROCESSING при стоящем NEW это воркеры, которые взяли и не вернули.
Просроченные аренды: сколько строк в PROCESSING с claimed_until < now() подобрали заново. Ненулевой фон говорит, что воркеры падают или аренда короче реальной обработки.
Длительность обработки одной задачи гистограммой с меткой типа задачи: по ней подбирают и размер пачки, и длину аренды.
На графике дежурного из этого стоят возраст и FAILED; остальное смотрят при разборе. И запись с причиной отказа в самой строке (last_error, обрезанная до пары сотен символов) экономит час поиска по логам.
Коротко
- Простой цикл на celery beat ломается при параллельных репликах, зависших записях и последовательных HTTP-вызовах без таймаутов.
- Правильный захват задач — атомарный UPDATE с
FOR UPDATE SKIP LOCKED: база сама распределяет работу между воркерами. claimed_until(«аренда задачи») защищает от зависших записей после перезапуска — просроченную аренду воркеры подберут автоматически.- HTTP-вызовы — между транзакциями, не внутри; обработка — параллельная через пул воркеров, не последовательная.
- Повторные попытки с нарастающей задержкой и терминальный статус
FAILEDне дают очереди застрять. - Фреймворк пакетной обработки — для конечных задач с перезапуском; брокер — для масштаба и нескольких потребителей.
- «Ровно один раз» по расписанию даёт
pg_try_advisory_lock, уникальная строка прогона илиCronJobсconcurrencyPolicy: Forbid; очередь наблюдают по возрасту старейшей задачи, статусам, просроченным арендам и длительности обработки.
Что почитать дальше
- Паттерны отказоустойчивости на Python — таймауты, повторные попытки, автоматический выключатель для вызова внешних систем.
- Распределённые паттерны на Python — transactional outbox и идемпотентность при работе с брокером.
- Сессия и транзакции в SQLAlchemy — где проходят границы транзакции и почему сетевой вызов внутри них ломает пул.