Интеграции и API

Как проектировать background jobs: retry, heartbeat, failed state и защита от повторного выполнения

Фоновая задача кажется простой, пока всё работает.

Есть очередь.

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_DOCUMENT

payload:

{
  "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 UUID

Provider видит две разные операции.

Хорошо:

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
→
PAID

Job должна выполнять переход только из допустимого состояния.

Если повтор приходит после первого успеха:

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
→ retry
RATE_LIMITED
→ retry later
AUTH_FAILED
→ failed / alert
PAYLOAD_INVALID
→ failed immediately
PROVIDER_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 = 2

Worker умеет:

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_at

Attempts

Распределение количества попыток.

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 sec
PDF report:
< 5 min
nightly 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 workers
heavy-processing workers
integration 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 ms

heartbeat каждые десять секунд бессмысленен.

Достаточно общего execution timeout.

Heartbeat особенно полезен для:

долгих jobs;
непредсказуемого времени;
jobs с checkpoints;

Для коротких задач проще

Например:

send webhook

обычно занимает:

< 5 sec.

Можно использовать:

lease = 60 sec

без регулярного heartbeat.

Если worker исчез, задача автоматически recover через минуту.


Не нужно одинаково усложнять все типы jobs

Job type может иметь свою policy.

Например:

JobTimeoutHeartbeatAttempts
SEND_EMAIL30 secНет5
WEBHOOK20 secНет8
GENERATE_PDF10 minДа3
IMPORT_DATA2 hДа + checkpoint3
DELETE_STORAGE60 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 = 0

Worker A claim:

status = running
attempts = 1
execution_id = abc
heartbeat = 10:00

CRM отвечает:

503

Job:

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:1842

CRM отвечает:

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 перестаёт зависеть от удачи.

И фоновая обработка становится тем, чем она должна быть в зрелом веб-продукте:

восстанавливаемым, наблюдаемым и предсказуемым механизмом выполнения работы после пользовательского запроса.

Есть похожая задача?

Опишите продукт, интеграции и ограничения. До разработки зафиксируем объём, риски и критерии приёмки.

Обсудить проект →