Фоновая задача кажется простой, пока всё работает.
Есть очередь.
Worker получил job.
Выполнил.
Поставил:
completedГотово.
Но production начинается в тот момент, когда между двумя этими действиями происходит что-нибудь неприятное.
Например:
worker получил задачу;отправил запрос во внешний API;внешний API успешно выполнил операцию;соединение оборвалось до получения ответа;worker решил, что задача не выполнена;очередь запустила её повторно.Если операция была:
отправить письмопользователь получил два письма.
Если:
создать заказпоявилось два заказа.
Если:
вернуть деньгипоследствия уже намного серьёзнее.
Именно поэтому production-ready background jobs — это не просто:
queue + workerа полноценная модель отказов:
claim
↓
execute
↓
heartbeat
↓
success / failure
↓
retry
↓
recovery
↓
terminal failed stateИ поверх всего этого — защита от повторного выполнения.
Разберём, как такую систему проектировать.
Зачем вообще нужны background jobs
Допустим пользователь нажал:
Создать проект.
Сервер должен:
записать проект;
создать историю;
отправить email;
отправить Telegram;
сгенерировать PDF;
обновить CRM;
вызвать внешний API.Если выполнить всё внутри HTTP-request:
Browser
↓
API
↓
PostgreSQL
↓
SMTP
↓
Telegram
↓
PDF
↓
External API
↓
Responseвремя ответа становится зависимым от всех внешних компонентов.
Достаточно одному из них работать медленно — пользователь ждёт.
Ещё хуже, если:
Project created ✓
SMTP unavailable ✗а API возвращает:
500Пользователь не понимает:
Проект создался или нет?
Поэтому критичную транзакцию полезно отделять от вторичных действий.
Например:
HTTP request
↓
Create project
↓
Create background jobs
↓
COMMIT
↓
200 OKА дальше:
Worker
├─ email
├─ PDF
└─ integrationsПользовательский запрос заканчивается быстро.
Фоновая инфраструктура самостоятельно доводит вторичные операции до результата.
Но «выполнить позже» создаёт новый класс проблем
Теперь между созданием задачи и её выполнением могут пройти:
секунды;
минуты;
часы;За это время:
worker может упасть;
сервер может перезапуститься;
внешний сервис может быть недоступен;
приложение может обновиться;
данные могут измениться.Поэтому job нужно рассматривать не как функцию:
await sendEmail();а как самостоятельную сущность с жизненным циклом.
Минимальная state machine background job
Практичный базовый вариант:
PENDING
↓
RUNNING
↓
┌───────────┐
│ │
▼ ▼
COMPLETED PENDING
retry
│
▼
FAILEDТо есть задача всегда находится в понятном состоянии.
Например:
pendingожидает worker.
runningуже кем-то выполняется.
completedуспешно закончена.
failedавтоматические попытки закончились или ошибка признана невосстановимой.
Почему не стоит удалять job сразу после успеха
Можно сделать:
job выполнена
↓
DELETE FROM jobsНо тогда теряется полезная эксплуатационная информация.
Например через два часа пользователь сообщает:
Письмо так и не пришло.
Что происходило?
Без истории ответить трудно.
Поэтому хотя бы некоторое время полезно хранить:
status;
attempts;
started_at;
finished_at;
last_error;
worker_id.Успешные задачи можно удалять позже по retention policy.
Пример таблицы
Для PostgreSQL queue это может выглядеть примерно так:
CREATE TABLE background_jobs (
id bigserial PRIMARY KEY,
queue text NOT NULL DEFAULT 'default',
type text NOT NULL,
version integer NOT NULL DEFAULT 1,
payload jsonb NOT NULL,
status text NOT NULL DEFAULT 'pending',
priority integer NOT NULL DEFAULT 0,
available_at timestamptz NOT NULL DEFAULT now(),
attempts integer NOT NULL DEFAULT 0,
max_attempts integer NOT NULL DEFAULT 5,
worker_id text,
started_at timestamptz,
heartbeat_at timestamptz,
last_error_code text,
last_error text,
progress jsonb,
created_at timestamptz NOT NULL DEFAULT now(),
completed_at timestamptz,
failed_at timestamptz
);Не каждому проекту нужны все поля сразу.
Но такая схема хорошо показывает жизненный цикл задачи.
Job начинается с PENDING
Например пользователь загрузил документ.
Создаётся:
PROCESS_DOCUMENTpayload:
{
"fileId": 4812
}Запись:
status = pending
attempts = 0
available_at = now()Теперь задача существует независимо от HTTP-request.
Даже если API process сразу после commit перезапустится, worker позже найдёт запись.
Worker должен сначала безопасно «захватить» job
Если workers несколько, нельзя просто сделать:
SELECT *
FROM background_jobs
WHERE status = 'pending'
LIMIT 1;Два worker могут получить одну строку одновременно.
Для PostgreSQL queue можно использовать короткий claim через:
FOR UPDATE SKIP LOCKEDНапример:
WITH selected AS (
SELECT id
FROM background_jobs
WHERE status = 'pending'
AND available_at <= now()
ORDER BY priority DESC, available_at, id
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE background_jobs AS j
SET
status = 'running',
worker_id = $1,
started_at = now(),
heartbeat_at = now(),
attempts = attempts + 1
FROM selected
WHERE j.id = selected.id
RETURNING j.*;После этого transaction должна закончиться.
Выполнять job внутри открытой transaction — плохая идея
Не стоит делать:
BEGIN
↓
claim job
↓
HTTP request 40 sec
↓
PDF generation 2 min
↓
COMMITТак transaction живёт слишком долго.
Это увеличивает:
время locks;
нагрузку на MVCC;
риск конфликтов;
сложность VACUUM;и просто без необходимости связывает внешний мир с открытой database transaction.
Лучше:
CLAIM
↓
COMMIT
↓
EXECUTE
↓
короткая transaction
для фиксации результатаJob теперь RUNNING
Worker получил:
job #4812и начинает работу.
Самый простой happy path:
RUNNING
↓
execute
↓
success
↓
COMPLETEDНапример:
UPDATE background_jobs
SET
status = 'completed',
completed_at = now(),
worker_id = NULL,
heartbeat_at = NULL
WHERE id = $1;Но production редко ограничивается happy path.
Первая проблема: временная ошибка
Worker обращается к внешнему API.
Ответ:
503 Service UnavailableЭто не обязательно означает:
Задача никогда не выполнится.
Вероятно, внешний сервис восстановится.
Поэтому job нужно вернуть:
RUNNING
↓
temporary failure
↓
PENDINGс новой датой:
available_at = future.Это retry.
Не каждая ошибка должна повторяться
Очень важно разделять ошибки хотя бы на несколько типов.
Временная
Например:
network timeout;
HTTP 429;
HTTP 503;
temporary database error.Логично retry.
Постоянная
Например:
email address invalid;
file format unsupported;
required entity permanently deleted.Повторить через минуту ничего не изменит.
Неизвестная
Например неожиданный exception.
Здесь часто лучше ограниченно повторить, но обязательно сохранить ошибку и наблюдать её.
Retry policy должна быть явной
Плохо:
catch (error) {
retry();
}для абсолютно любой ошибки.
Лучше иметь семантику:
if (error instanceof RetryableError) {
scheduleRetry();
}
if (error instanceof PermanentError) {
markFailed();
}И отдельную обработку неизвестных ошибок.
Почему мгновенный retry опасен
Представим внешний API упал.
В очереди:
20 000 jobs.Все workers начинают:
request
→ error
→ retry
→ request
→ error
→ retryЧерез несколько секунд собственная система начинает атаковать неисправный сервис тысячами запросов.
И заодно тратит:
CPU;
network;
database writes;
logs.Retry должен давать зависимости время восстановиться.
Нужен backoff
Например:
attempt 1
→ 30 sec
attempt 2
→ 2 min
attempt 3
→ 10 min
attempt 4
→ 30 min
attempt 5
→ failedМожно использовать exponential backoff:
delay =
base × 2^(attempt - 1)с верхним пределом.
И желательно добавлять jitter
Представим 10 000 jobs получили ошибку ровно в:
12:00:00Если всем поставить:
retry at 12:01:00через минуту система создаст новый пик.
Лучше слегка разнести попытки:
12:00:52
12:00:57
12:01:04
12:01:11
...Это и есть jitter.
Для rate limit нужен особый retry
Например provider возвращает:
429 Too Many Requestsи:
Retry-After: 120Нет смысла игнорировать это и использовать универсальный:
retry через 5 секунд.Если внешний сервис сообщает окно восстановления, worker может использовать его.
Второй важный механизм — heartbeat
Представим job обрабатывает большой файл.
Обычное время:
15 минут.Worker поставил:
status = runningи начал обработку.
Через две минуты process погиб:
OOMВ базе всё ещё:
running.Кто должен понять, что задача больше никем не выполняется?
Для этого нужен lease
Когда worker забирает job, он фактически говорит:
Я владею этой задачей, пока подтверждаю, что жив.
Самый простой вариант:
started_atи фиксированный timeout.
Например:
job timeout = 10 minЕсли:
now - started_at > 10 minсчитаем worker потерянным.
Но фиксированного started_at недостаточно для долгих задач
Job может совершенно нормально работать:
45 минут.Поэтому появляется:
heartbeat_at.Worker периодически обновляет:
UPDATE background_jobs
SET heartbeat_at = now()
WHERE id = $1
AND worker_id = $2
AND status = 'running';Например каждые:
30 секунд.Recovery проверяет свежесть heartbeat
Например:
heartbeat timeout = 2 min.Job:
status = running
heartbeat_at =
7 минут назадочевидно подозрительна.
Recovery process может решить:
worker потерян.Внешние системы используют похожую идею
Например в queue systems существует concept visibility timeout: consumer получает сообщение, и оно временно становится невидимым для других consumers; если обработка затягивается, lease можно продлевать. Для долгой работы это фактически родственник heartbeat-модели.
Для собственной job system смысл тот же:
worker должен периодически подтверждать право продолжать обработку.
Heartbeat особенно полезен для задач с progress
Например импорт:
10 000 записей.Worker обработал:
3 700.Можно хранить:
{
"processed": 3700,
"total": 10000
}вместе с heartbeat.
Администратор видит:
37%а мониторинг понимает:
Job действительно движется.
Но progress и heartbeat — разные вещи
Задача может быть живой, но долго находиться на одном этапе.
Например:
uploading large object.Поэтому heartbeat должен означать:
Worker жив и всё ещё контролирует execution.
А progress:
Вот сколько полезной работы уже выполнено.
Не следует требовать изменения progress при каждом heartbeat.
Recovery после потерянного heartbeat
Упрощённо:
RUNNING
↓
heartbeat stale
↓
worker considered dead
↓
job recoveredДальше возможны два варианта.
Если:
attempts < max_attemptsвозвращаем:
PENDING.Если лимит достигнут:
FAILED.Но вот здесь появляется главная проблема background jobs
Представим worker сделал:
отправил emailпосле чего process мгновенно упал.
Он не успел поставить:
COMPLETED.Через heartbeat timeout recovery делает:
PENDING.Другой worker выполняет job снова.
Пользователь получает:
два письма.Повторное выполнение — не исключение
Это принципиально важная мысль.
Production job system следует проектировать так, будто job может быть выполнена больше одного раза.
Не:
Ну это случится только при очень редкой аварии.
А:
Это нормальный failure mode распределённой системы.
Потому что всегда существует промежуток между:
side effect completedи:
queue knows it completed.В этом промежутке может произойти сбой.
Поэтому нужна идемпотентность
Идемпотентная операция допускает повторный запрос без повторного бизнес-эффекта.
Например job:
ACTIVATE_SUBSCRIPTION
user = 52
plan = PROПервое выполнение:
ACTIVE.Второе:
уже ACTIVE
→ ничего нового не создаём.Это естественная идемпотентность.
Но SEND_EMAIL сам по себе не идемпотентен
Если просто вызвать:
sendEmail();два раза — будет два письма.
Значит нужен дополнительный механизм.
Вариант 1: business idempotency key
Например:
PROJECT_CREATED_EMAIL:project:1842В таблице side effects:
CREATE TABLE processed_effects (
key text PRIMARY KEY,
created_at timestamptz NOT NULL
);Worker сначала пытается:
INSERT INTO processed_effects(key)
VALUES ('PROJECT_CREATED_EMAIL:project:1842');Если:
INSERT successон первый.
Если:
unique violationтакая операция уже запускалась.
Но здесь есть тонкость
Если записать idempotency key:
до emailа затем email не отправится, повтор будет ошибочно заблокирован.
Если записать:
после emailмежду email и INSERT опять есть окно crash.
То есть внешние side effects нельзя сделать exactly-once одной локальной таблицей магически.
Намного лучше, если внешний API сам поддерживает idempotency
Например job создаёт объект у payment provider.
Отправляем стабильный ключ:
refund:payment-1842:fullПервая попытка:
provider creates refund.Connection потерялась.
Вторая попытка отправляет тот же ключ.
Provider понимает:
Это повтор той же операции.
и не создаёт второй refund.
Именно так idempotency key превращает сетевой retry из опасной операции в ожидаемый сценарий.
Ключ должен описывать бизнес-операцию, а не попытку
Плохо:
attempt-1: random UUID
attempt-2: another UUIDProvider видит две разные операции.
Хорошо:
refund:payment-1842:fullдля всех retries одной операции.
Вариант 2: database constraint как последняя защита
Допустим job создаёт invoice.
Бизнес-правило:
Для заказа может существовать только один invoice этого типа.
Добавляем:
UNIQUE(order_id, invoice_type)Теперь даже если job выполнится параллельно два раза:
invoice #1 ✓
invoice #2 → UNIQUE violationБаза защищает invariant.
Constraints часто надёжнее application-lock
Можно использовать:
Redis lock;
advisory lock;
mutex.Но process может упасть.
Lease может истечь.
Network partition может нарушить предположения.
Если invariant можно выразить в PostgreSQL:
UNIQUE
CHECK
FOREIGN KEYэто очень сильная последняя линия защиты.
Вариант 3: compare-and-set состояния
Job:
CONFIRM_ORDERможет выполнить:
UPDATE orders
SET status = 'confirmed'
WHERE id = $1
AND status = 'pending';Если:
affected rows = 1переход произошёл.
Если:
0возможно заказ уже подтверждён.
Повторная job ничего не меняет.
Это особенно удобно для state machine
Например:
WAITING_PAYMENT
→
PAIDJob должна выполнять переход только из допустимого состояния.
Если повтор приходит после первого успеха:
current state = PAIDоперация становится no-op.
Повторное выполнение можно сделать частью domain model
Это намного лучше, чем:
попробуем гарантировать,
что worker никогда не запустит job дважды.Такой гарантии в реальной распределённой системе добиться значительно сложнее.
Exactly-once часто является неправильной целью
На уровне бизнеса обычно полезнее сформулировать:
Физическая обработка может происходить несколько раз, но бизнес-эффект должен появиться один раз.
Например worker может три раза получить:
CREATE_INVOICEНо в базе существует:
1 invoice.Вот это действительно полезная гарантия.
Третий важный элемент — FAILED state
Очень плохая система retries выглядит так:
ошибка
↓
retry
↓
ошибка
↓
retry
↓
ошибка
↓
retry
↓
...Job никогда не исчезает.
Но никогда и не завершается.
Нужна terminal state
Например:
attempt 1 → fail
attempt 2 → fail
attempt 3 → fail
attempt 4 → fail
attempt 5 → FAILEDТеперь система признаёт:
Автоматически эту задачу решить не удалось.
Это важное состояние.
FAILED — не то же самое, что DELETE
Ошибка не должна исчезать.
Полезно сохранить:
job type;
payload;
attempts;
последнюю ошибку;
error code;
время;
worker;Например:
FAILED
type:
SYNC_CRM
attempts:
5
error:
CRM_TOKEN_REVOKEDАдминистратор сразу понимает:
Ещё один retry через минуту бессмысленен.
Нужно обновить credentials.
FAILED state — аналог dead-letter queue
В специализированных очередях часто используется отдельная dead-letter queue.
В собственной PostgreSQL job system можно добиться похожей операционной модели просто терминальным:
status = failed.И отдельным административным интерфейсом:
Failed jobsЧто должен уметь администратор
Например:
просмотреть ошибку;
открыть payload;
увидеть attempts;
увидеть последний heartbeat;
повторить задачу вручную;
отменить;Но manual retry тоже должен соблюдать idempotency.
Кнопка:
Повторитьне отменяет правила повторного выполнения.
Не стоит просто менять FAILED обратно на PENDING
Иногда лучше создавать:
новую execution attemptили хотя бы записывать:
manual_retry_count;
retried_by;
retried_at.Так сохраняется audit trail.
Почему error code лучше одного текста
Плохо:
last_error =
"Something went wrong"Гораздо полезнее:
error_code =
CRM_AUTH_FAILED
last_error =
HTTP 401 from CRM providerТеперь monitoring может агрегировать:
CRM_AUTH_FAILED = 417 jobsи сразу заметить общий инцидент.
Errors полезно классифицировать централизованно
Например:
NETWORK_TIMEOUT
RATE_LIMITED
AUTH_FAILED
PAYLOAD_INVALID
ENTITY_NOT_FOUND
PROVIDER_5XX
UNKNOWNИ retry policy может зависеть от класса.
Например
NETWORK_TIMEOUT
→ retryRATE_LIMITED
→ retry laterAUTH_FAILED
→ failed / alertPAYLOAD_INVALID
→ failed immediatelyPROVIDER_5XX
→ retryТакая система намного предсказуемее.
Что делать, если пользователь удалил объект, пока job ждала
Например очередь содержит:
GENERATE_PROJECT_REPORT
projectId = 1842Но пользователь уже удалил или архивировал проект.
Worker открывает job через час.
Что теперь?
Ответ определяется domain.
Иногда правильный результат — success/no-op
Если отчёт больше не нужен:
project missing
→ job no longer applicable
→ completed/skipped.Нет смысла retry пять раз.
Иногда это ошибка
Если job:
CAPTURE_PAYMENTа order внезапно отсутствует — это может означать повреждение состояния.
Тогда:
FAILED
+
alert.Следовательно, status skipped иногда тоже полезен
Например:
PENDING
↓
business condition no longer relevant
↓
SKIPPEDОн отличается от:
FAILED.Система не сломалась.
Просто операция потеряла смысл.
Payload должен быть маленьким и стабильным
Есть соблазн положить в job:
{
"project": {
"...": "вся текущая сущность"
},
"user": {
"...": "весь пользователь"
}
}Но через день данные изменятся.
Какая версия правильная?
Чаще лучше хранить ID
Например:
{
"projectId": 1842,
"userId": 52
}Worker загружает актуальное состояние.
Но иногда нужен snapshot
Например:
Отправить пользователю текст письма именно в том виде, в котором он был сформирован в момент события.
Если template завтра изменится, retry не должен вдруг отправить другое письмо.
Тогда можно сохранить:
templateVersion;
rendering data;или уже подготовленный immutable snapshot.
Главное — решить это осознанно
Вопрос:
Job использует данные на момент создания или на момент исполнения?
должен иметь ответ.
Delayed jobs особенно чувствительны к versioning
Представим job:
SEND_REMINDERсоздана сегодня.
available_at:
через 30 дней.За месяц приложение обновилось десять раз.
Payload v1 больше не соответствует worker v8.
Поэтому job contract тоже может иметь version
Например:
type = SEND_REMINDER
version = 2Worker умеет:
v1 → старый decoder
v2 → новый decoderили существует контролируемая migration pending jobs.
Нельзя считать очередь временной памятью
Некоторые jobs могут жить достаточно долго.
Значит они являются persistent data contract между версиями приложения.
Это особенно важно при rolling deployment, когда одновременно работают:
Worker v7
Worker v8.Graceful shutdown worker
Представим deploy.
Операционная система отправляет:
SIGTERM.Worker в этот момент выполняет:
EXPORT_DATAуже 8 минут.
Плохое поведение:
SIGTERM
↓
process exits immediately.Теперь recovery запустит job снова.
Более аккуратная схема
После SIGTERM:
stop claiming new jobs
↓
continue current job
↓
finish / checkpoint
↓
exitс ограниченным:
grace period.Почему grace period всё равно нужен
Job может зависнуть навсегда.
Нельзя останавливать deployment:
до бесконечности.Поэтому:
SIGTERM
↓
30/60/... seconds grace
↓
force terminationа дальше job восстанавливается через heartbeat/lease.
Для долгих jobs полезны checkpoints
Например импорт:
1 000 000 строк.Worker дошёл до:
row 700 000и сервер перезапустился.
Начинать всё с нуля дорого.
Job может сохранять checkpoint
Например:
{
"lastProcessedId": 700000
}При retry:
resume from checkpoint.Heartbeat может одновременно обновлять checkpoint.
Но checkpoint тоже должен быть корректным
Нельзя записать:
processed = 700000до того, как первые 700 000 действительно durable.
Иначе после crash job пропустит часть данных.
Хороший паттерн — checkpoint после завершённого batch
Например:
process 1000 records
↓
commit result
↓
save checkpoint
↓
next batchТеперь повтор максимум переобработает безопасный небольшой участок.
Background job и workflow — не всегда одно и то же
Job хорошо подходит для:
одного ограниченного действия.Например:
send email;
delete object;
generate PDF;
call webhook.Но процесс:
создать счёт;
ждать оплату 7 дней;
при оплате создать доступ;
при отсутствии оплаты отправить напоминание;
ещё через 3 дня отменить.уже не очень похож на одну job.
Это скорее workflow
Состояние живёт долго.
Есть:
несколько шагов;
ожидания;
внешние события;
таймеры;
компенсации.Можно построить workflow поверх jobs.
Но не стоит превращать одну job в:
гигантскую функцию,
которая неделю пытается жить.Хорошая job ограничена по ответственности
Например:
SEND_INVOICE_EMAILлучше, чем:
HANDLE_EVERYTHING_AFTER_ORDERкоторая:
создаёт invoice;
отправляет письмо;
обновляет CRM;
создаёт документы;
шлёт push;
синхронизирует аналитику.Если один этап падает, становится непонятно:
Какие предыдущие уже выполнены?
Лучше отдельные durable effects
Например:
ORDER_PAID
│
├── ISSUE_DOCUMENT
├── SEND_EMAIL
├── CRM_SYNC
└── ANALYTICS_EVENTКаждая задача:
retry;
failed state;
monitoringимеет отдельно.
Но fan-out тоже должен быть атомарным
Если все jobs создаются после бизнес-транзакции в том же PostgreSQL, можно вставить их одной transaction.
Например:
BEGIN;
UPDATE orders
SET status = 'paid'
WHERE id = $1;
INSERT INTO background_jobs (...)
VALUES (...);
INSERT INTO background_jobs (...)
VALUES (...);
INSERT INTO background_jobs (...)
VALUES (...);
COMMIT;Так бизнес-факт и необходимая фоновая работа появляются согласованно.
При внешнем broker нужен transactional outbox
Если задачи публикуются в RabbitMQ/Kafka после commit, возникает dual-write.
Один из стандартных способов решения:
business transaction
↓
outbox row
↓
COMMIT
↓
publisher
↓
brokerТо есть даже при специализированном broker PostgreSQL часто всё равно участвует в гарантированной публикации события.
Ограничение concurrency
Представим 20 workers.
Все получили jobs для:
customer 1842.И одновременно вызывают API одного партнёра.
Или пытаются изменить один и тот же файл.
Общий worker concurrency = 20 не означает, что для каждой сущности допустим concurrency = 20.
Иногда нужен per-key serialization
Например:
одновременно только 1 job
на один payment.Или:
не больше 2 jobs
на одного external provider account.Это можно реализовать:
database locks;
advisory locks;
dedicated queues;
partitioning по ключу.Конкретный механизм зависит от задачи.
Но lock не заменяет idempotency
Даже если сейчас два worker не выполнят операцию одновременно, job всё равно может быть повторена позже после crash.
Поэтому:
mutual exclusionи:
idempotencyрешают разные проблемы.
Timeout должен существовать на нескольких уровнях
Например job вызывает внешний HTTP API.
Нужны:
connection timeout;
request timeout;
job execution timeout;
heartbeat timeout.Это разные значения.
Без request timeout worker может зависнуть:
навсегда.Heartbeat при этом тоже может прекратиться, и другой worker начнёт job повторно, пока первый process фактически ещё жив, но висит в сетевом вызове.
Значит worker должен уметь отменять зависшее действие
Когда job execution deadline превышен:
AbortControllerили соответствующий механизм должен попытаться остановить внешний request.
После этого job переходит в понятный retry/failure path.
Heartbeat не должен скрывать вечную работу
Плохая job:
while (true) {
heartbeat();
}формально всегда здорова.
Но результата нет.
Поэтому отдельно полезен:
maximum execution durationили domain-specific deadline.
Heartbeat отвечает:
Worker жив.
Timeout:
Но слишком долго жить этой job всё равно нельзя.
Мониторинг background jobs
Просто видеть:
Worker process UPнедостаточно.
Нужно наблюдать сам поток работы.
Минимальные метрики
Pending
jobs_pendingСколько задач ожидает.
Running
jobs_runningСколько сейчас выполняется.
Failed
jobs_failedСколько дошло до terminal failure.
Oldest pending age
now - oldest available_atОдна из самых полезных метрик.
Processing duration
completed_at - started_atAttempts
Распределение количества попыток.
Recovered jobs
Сколько running пришлось возвращать после stale heartbeat.
Почему oldest pending age важнее простого COUNT
Например:
pending = 10 000но processing rate:
20 000 jobs/sec.Очередь практически пуста.
А может быть:
pending = 4но:
oldest = 6 hours.Четыре задачи явно застряли.
Количество само по себе мало что говорит.
Полезен throughput
Например:
jobs completed / minuteВ обычный день:
500/min.Сегодня:
7/min.Workers существуют.
Heartbeat приходит.
Но производительность обрушилась.
Это operational incident.
Ошибки нужно агрегировать по коду
Например:
CRM_AUTH_FAILED = 917за пять минут.
Гораздо информативнее:
917 random error strings.Можно сразу понять:
У CRM-интеграции истёк token.
Alert нужен не на каждую failed job
Некоторые ошибки ожидаемы.
Например пользователь указал:
недействительный email.Это не повод будить инженера ночью.
Лучше алертить тенденции
Например:
failed rate > threshold;oldest pending > SLA;stale running jobs > 0 длительное время;worker heartbeat absent;queue depth постоянно растёт.Background jobs должны иметь SLA
Например:
password email:
< 30 secPDF report:
< 5 minnightly cleanup:
до утраТогда monitoring знает, что значит:
слишком долго.Без SLA queue age является просто цифрой.
Priorities
Не все задачи одинаково срочные.
Например:
PASSWORD_RESET_EMAILнельзя ставить за:
100 000 analytics jobs.Можно использовать:
priority.Но осторожно.
Постоянно высокий priority может вызвать starvation
Если критичные jobs появляются непрерывно:
low priorityникогда не выполняется.
Можно использовать:
отдельные queues;
reserved worker capacity;
aging priority.Иногда отдельные worker pools проще
Например:
email workersheavy-processing workersintegration workersТеперь тяжёлая обработка PDF не блокирует password reset.
Отмена job
Пользователь запустил экспорт.
Через минуту нажал:
Отменить.
Job уже RUNNING.
Просто поставить:
status = cancelledнедостаточно, если worker продолжает выполнять работу.
Cancellation должна быть cooperative
Например worker между этапами проверяет:
cancel_requested_at.Если установлено:
cleanup
↓
stop
↓
CANCELLEDДля внешнего request может потребоваться abort.
Не каждую операцию вообще можно отменить
Если payment provider уже выполнил refund, кнопка:
Отменить jobне вернёт деньги обратно автоматически.
То есть job cancellation и business compensation — разные понятия.
Security payload
Queue часто содержит:
email;
user IDs;
document IDs;
webhook data.Не стоит складывать туда:
пароли;
access tokens;
полные банковские данные;если этого можно избежать.
Особенно опасен last_error
Некоторые HTTP libraries записывают в exception:
request headers;
authorization token;
full response.Если просто сделать:
last_error = error.stackв административной таблице могут появиться secrets.
Ошибка должна проходить sanitization.
Payload retention тоже важен
Даже completed jobs могут содержать персональные данные.
Если для диагностики достаточно хранить successful jobs:
14 дней,не обязательно оставлять их навсегда.
FAILED jobs можно хранить дольше, но тоже в рамках осознанной политики.
Backups и queue
Если jobs находятся в PostgreSQL, они автоматически могут попасть в database backup.
Это полезно.
Но нужно понимать последствия restore.
Представим восстановили БД вчерашней давности
В ней есть:
PENDING jobsкоторые в реальном прошлом уже были выполнены.
После restore workers запустятся и выполнят их снова.
Это ещё одна причина обязательной идемпотентности
Disaster recovery способен создавать повторную обработку даже без crash конкретного worker.
После restore состояние очереди и внешнего мира может не совпадать.
Перед запуском workers после большого restore полезен reconciliation
Например:
payments;
emails;
webhooks;
files.Нужно понять, какие side effects уже произошли вне восстановленной БД.
Особенно для необратимых операций.
Testing background jobs
Обычный unit test:
execute(job)
→ successслишком слабый.
Production-система должна тестировать сбои.
Тест №1. Обычный успех
PENDING
↓
RUNNING
↓
COMPLETEDПроверяем:
attempts = 1;
completed_at set;Тест №2. Временная ошибка
attempt 1
→ provider 503Ожидаем:
PENDING
available_at > now
attempts = 1Тест №3. Permanent error
INVALID_PAYLOADОжидаем:
FAILEDбез бессмысленных retries.
Тест №4. Exhausted retries
attempt 1 fail
attempt 2 fail
attempt 3 failПосле max_attempts:
FAILED.Тест №5. Worker погиб
claim job
↓
RUNNING
↓
kill -9 workerЖдём heartbeat timeout.
Другой worker должен:
recover
↓
retry.Тест №6. Crash после успешного side effect
Это один из самых ценных тестов.
external API
→ SUCCESS
↓
kill worker
до markCompleted()Job повторяется.
Проверяем:
Создался ли второй бизнес-эффект?
Если да — job ещё не production-safe.
Тест №7. Retry внешнего API с одинаковым idempotency key
Первая попытка:
timeout after provider success.Вторая отправляет тот же:
Idempotency-Key.Проверяем, что внешний объект один.
Тест №8. Graceful deployment
Worker выполняет job.
Отправляем:
SIGTERM.Ожидаем:
new jobs no longer claimed;
current finishes;
process exits.Тест №9. Poison job
Job всегда вызывает:
PAYLOAD_INVALID.Она не должна бесконечно циркулировать.
Ожидаем:
FAILED.Тест №10. Backlog
Создаём:
100 000 jobs.Проверяем:
claim latency;
PostgreSQL load;
worker throughput;
oldest age.Production queue должна проходить не только correctness, но и load scenario.
Тест №11. Два worker одновременно
Оба пытаются получить одну job.
Проверяем корректную claim semantics.
Для side effect всё равно оставляем idempotency, потому что race во время claim — не единственный источник повторов.
Тест №12. PostgreSQL restart
Во время backlog перезапускаем базу.
После восстановления:
pending preserved;
workers reconnect;
queue continues;
stale running recovered.Тест №13. Release с pending jobs старой версии
Создаём:
payload version 1.Обновляем worker до новой версии.
Он должен:
обработать v1или явно выполнить предусмотренную migration.
Не:
JSON parse error → failed 20 000 jobs.Worker loop в упрощённом виде
Концептуально:
while (!shuttingDown) {
const job = await claimNextJob();
if (!job) {
await waitForJobs();
continue;
}
try {
startHeartbeat(job);
await execute(job);
await markCompleted(job);
} catch (error) {
if (isPermanent(error)) {
await markFailed(job, error);
continue;
}
if (job.attempts >= job.maxAttempts) {
await markFailed(job, error);
continue;
}
await scheduleRetry(
job,
calculateBackoff(job.attempts)
);
} finally {
stopHeartbeat(job);
}
}Реальная production-реализация будет сложнее.
Но важен именно жизненный цикл.
Recovery loop
Отдельный процесс периодически ищет:
status = runningс устаревшим:
heartbeat_at.Дальше:
if attempts >= max_attempts
→ FAILED
else
→ PENDINGи очищает:
worker_id;
started_at;
heartbeat_at.Worker ownership тоже нужно проверять
Представим старая execution проснулась после долгого зависания.
А job уже recovered и выполняется новым worker.
Старый worker не должен затем сделать:
UPDATE job
SET status = completedповерх нового владельца.
Полезен execution token
При claim создаём:
execution_id = random UUID.Heartbeat:
UPDATE background_jobs
SET heartbeat_at = now()
WHERE id = $jobId
AND execution_id = $executionId
AND status = 'running';Completion:
UPDATE background_jobs
SET status = 'completed'
WHERE id = $jobId
AND execution_id = $executionId
AND status = 'running';Если job уже была recovered и получила новый execution ID, старый worker больше не может завершить новую execution.
Это защищает от zombie worker
Worker считался погибшим.
Lease истёк.
Задачу забрал другой.
Но старый process неожиданно разморозился.
Без ownership token два worker могут одновременно считать себя владельцами job.
Heartbeat и execution token работают вместе
execution_id
→ кто владелец
heartbeat_at
→ жив ли этот владелецПолучается намного более точная lease-модель.
Нужно ли делать heartbeat для всех jobs?
Не обязательно.
Если job обычно занимает:
50 msheartbeat каждые десять секунд бессмысленен.
Достаточно общего execution timeout.
Heartbeat особенно полезен для:
долгих jobs;
непредсказуемого времени;
jobs с checkpoints;Для коротких задач проще
Например:
send webhookобычно занимает:
< 5 sec.Можно использовать:
lease = 60 secбез регулярного heartbeat.
Если worker исчез, задача автоматически recover через минуту.
Не нужно одинаково усложнять все типы jobs
Job type может иметь свою policy.
Например:
| Job | Timeout | Heartbeat | Attempts |
|---|---|---|---|
SEND_EMAIL | 30 sec | Нет | 5 |
WEBHOOK | 20 sec | Нет | 8 |
GENERATE_PDF | 10 min | Да | 3 |
IMPORT_DATA | 2 h | Да + checkpoint | 3 |
DELETE_STORAGE | 60 sec | Нет | 10 |
Это намного разумнее одной глобальной настройки.
Retry тоже должен зависеть от типа задачи
Email:
5 attempts.Удаление мусора из storage:
10 attempts
в течение суток.Payment action:
ограниченный retry
+
обязательная idempotency.Некоторые jobs вообще:
max_attempts = 1если повтор опасен и внешняя система не предоставляет механизм безопасного retry.
Самая опасная job — необратимый side effect без idempotency
Например внешний сервис имеет API:
POST /chargeи не предоставляет:
idempotency key;
lookup by external request id;
deduplication.Worker отправил charge.
Connection потерялась.
Что делать?
Автоматический retry может быть опасен
Потому что:
первая операция могла пройти.В таких случаях job может перейти:
UNKNOWNили:
MANUAL_REVIEWвместо автоматического повторения.
Иногда нужен отдельный ambiguous state
Простой автомат:
PENDING
RUNNING
COMPLETED
FAILEDподходит не всем операциям.
Для финансовых/необратимых внешних действий бывает полезно:
PENDING_CONFIRMATIONили:
UNKNOWN.Это означает:
Мы не знаем, выполнил ли внешний provider операцию.
Сначала выполняем reconciliation.
Reconciliation важнее слепого retry
Например:
charge request timeout.Вместо:
отправить ещё разсистема может запросить:
GET transaction by business reference.Если transaction уже существует:
mark completed.Если нет:
retry create.Это намного безопаснее.
Background jobs должны моделировать неопределённость сети
Фраза:
Запрос завершился timeout.
не означает:
операция не произошла.Она означает:
мы не получили ответ.Это принципиальная разница.
Именно поэтому idempotency и reconciliation — часть worker architecture
А не дополнительные улучшения «на потом».
Когда job system становится workflow engine
Постепенно можно добавить:
retry;
heartbeat;
timers;
child jobs;
signals;
cancellation;
checkpoints;
compensation.И однажды обнаружить, что команда практически написала собственный workflow engine.
Это нормальный момент для переоценки
Если большинство процессов:
многошаговые;
живут днями;
ждут внешних событий;
требуют сложных компенсаций;специализированный workflow engine может оказаться экономически разумнее дальнейшего расширения собственной jobs-системы.
Но для обычного веб-продукта полноценная job infrastructure может оставаться простой
Например:
PostgreSQL
+
SKIP LOCKED
+
Workers
+
Retry
+
Heartbeat для длинных задач
+
Failed state
+
Idempotency
+
Monitoringэтого достаточно для очень большого количества production-сценариев.
Пример полной жизни job
Допустим нужно синхронизировать проект с CRM.
Создание:
status = pending
attempts = 0Worker A claim:
status = running
attempts = 1
execution_id = abc
heartbeat = 10:00CRM отвечает:
503Job:
status = pending
available_at = 10:01
last_error_code = CRM_TEMPORARY_UNAVAILABLEВторая попытка:
10:01
Worker B
attempts = 2
execution_id = defНачинает синхронизацию.
CRM создаёт запись.
Worker B погибает до completion.
Через lease timeout:
heartbeat stale.Recovery:
status = pending.Третья попытка:
Worker C отправляет тот же business idempotency key:
crm-project:1842CRM отвечает:
already processedили возвращает существующий объект.
Worker C:
status = completed.С точки зрения физической обработки job выполнялась три раза.
С точки зрения бизнеса CRM-проект появился один раз.
Вот это и есть качественная модель.
Что не следует делать
Не считать RUNNING гарантией единственного исполнения
Это всего лишь lease.
Не retry абсолютно всё
Ошибки бывают permanent.
Не делать retry без backoff
Иначе failure превращается в storm.
Не использовать один timeout для всех jobs
PDF и email живут в разных временных масштабах.
Не удалять failed jobs молча
Они нужны для диагностики.
Не считать heartbeat доказательством прогресса
Worker может быть жив, но застрять логически.
Не рассчитывать на «ровно один запуск»
Проектируйте idempotent business effect.
Не хранить secrets в payload/error logs
Очередь — часть persistent infrastructure.
Не считать процесс worker UP доказательством работы очереди
Нужны business metrics.
Практический production checklist background jobs
Перед запуском мы бы проверили:
- у каждой job есть явный state;
- claim атомарен и конкурентно безопасен;
- выполнение происходит вне длинной database transaction;
- job имеет execution timeout;
- для длинных задач предусмотрен heartbeat;
- recovery умеет находить stale
running; - stale execution не может завершить job после передачи новому worker;
- retry выполняется только для retryable errors;
- используется backoff;
- при массовых ошибках применяется jitter;
- существует
max_attempts; - после исчерпания попыток job попадает в наблюдаемый
failed; - failed jobs не удаляются без следа;
- у ошибок есть структурированные codes;
- manual retry сохраняет audit;
- side effects спроектированы идемпотентно;
- внешние idempotency keys стабильны между retries;
- database constraints защищают ключевые business invariants;
- ambiguous external result не retry вслепую, если это опасно;
- для неопределённых операций предусмотрен reconciliation;
- job payload имеет понятный contract;
- delayed jobs совместимы с будущими версиями worker;
- graceful shutdown прекращает claim новых задач;
- heartbeat/progress не записываются чрезмерно часто;
- существует retention completed/failed jobs;
- monitoring знает pending, running, failed и oldest age;
- отслеживаются processing duration и retry rate;
- workers имеют heartbeat/health;
- очереди разных классов нагрузки не блокируют друг друга;
- security-аудит проверяет payload и error logs на secrets;
- recovery после database restore учитывает возможное повторное выполнение;
- тесты действительно убивают worker посреди job и проверяют восстановление.
Если эти свойства реализованы, background processing становится настоящей production-подсистемой.
Что получает бизнес
На первый взгляд всё это выглядит исключительно технической работой.
Но результат очень конкретный.
Вместо:
Письмо иногда не отправляется.
получаем:
job failed;
attempts = 5;
reason = SMTP_AUTH_FAILED.Вместо:
Иногда создаются двойные операции.
получаем:
stable idempotency key;
unique invariant;
safe retry.Вместо:
Worker упал, и непонятно, что потерялось.
получаем:
stale heartbeat
↓
automatic recovery.Вместо:
Система уже три дня повторяет одну сломанную задачу.
получаем:
max attempts
↓
FAILED
↓
alert.Хорошая job system делает ошибки видимыми
Это, пожалуй, одно из её главных свойств.
Production-система не обязана никогда ошибаться.
Внешний API будет недоступен.
Сервер когда-нибудь перезапустится.
Network timeout обязательно случится.
Worker когда-нибудь погибнет ровно между двумя важными строками кода.
Качественная архитектура отличается тем, что после этого система не оказывается в неизвестном состоянии.
Она умеет сказать:
эта задача ожидает;эту сейчас выполняют;эта будет повторена;эта потеряла worker;эта восстановлена;эта выполнена;эта окончательно не выполнилась.Вместо вывода
Background jobs часто начинают с простой функции:
queue.add(job);Но production-надежность появляется не в момент, когда задача попала в очередь.
Она появляется, когда система знает, что делать во всех промежуточных состояниях.
Правильная модель выглядит примерно так:
CREATE
↓
PENDING
↓
CLAIM
↓
RUNNING
│
├── heartbeat
├── progress
│
├── SUCCESS
│ ↓
│ COMPLETED
│
└── ERROR
↓
retryable?
/ \
yes no
│ │
BACKOFF FAILED
│
PENDINGА поверх неё существует ещё одно обязательное правило:
JOB MAY RUN AGAIN.Именно поэтому worker нельзя проектировать так, будто каждую задачу он увидит ровно один раз.
Безопаснее исходить из обратного:
любая значимая job однажды будет повторена в самый неудобный момент.
После crash.
После timeout.
После restore.
После потери acknowledgement.
После ручного retry.
И если бизнес-операция остаётся корректной даже тогда, очередь действительно готова к production.
Поэтому четыре базовых механизма — retry, heartbeat, failed state и idempotency — на самом деле решают четыре разных вопроса:
Retry:
можно ли попробовать ещё раз?Heartbeat:
жив ли текущий исполнитель?Failed:
когда прекратить автоматические попытки
и показать проблему человеку?Idempotency:
что произойдёт,
если задача всё-таки выполнится повторно?Когда ответы на эти вопросы формализованы в коде, состояние job перестаёт зависеть от удачи.
И фоновая обработка становится тем, чем она должна быть в зрелом веб-продукте:
восстанавливаемым, наблюдаемым и предсказуемым механизмом выполнения работы после пользовательского запроса.