open atlas
← Все проекты

backend · advanced · 6d

Планировщик задач

Планировщик задач cron + backoff с доставкой at-least-once, идемпотентными обработчиками и visibility timeout — чтобы ни одна задача не терялась молча, даже при краше воркера на середине выполнения.

Собирая планировщик задач с нуля, ты сталкиваешься со всеми ловушками распределённых систем в ограниченном контексте: что происходит при краше воркера на середине задачи, как возникают дублирующиеся триггеры из-за clock skew, почему экспоненциальный backoff нуждается в jitter и как доказать идемпотентность. Postgres как очередь делает всё состояние проверяемым через SELECT — никакого брокера-чёрного ящика не нужно.

Результат

Планировщик, где задачи по cron выполняются at-least-once, упавшие задачи повторяются с экспоненциальным backoff, а chaos-тест с убийством воркеров на середине задачи даёт ноль потерь и ноль небезопасных дублей.

Этапы

0/6 · 0%
  1. 01Очерти очередь: модель данных, индекс по времени, долговечная постановка

    До любого воркер-цикла смоделируй очередь как таблицу Postgres и реши, что значит «долговечность». Каждая задача несёт состояние (pending → claimed → done → failed → dead), метку run_at, счётчик попыток и payload. Весь смысл Postgres-как-очереди в том, что состояние проверяемо через SELECT — без брокера-чёрного ящика, — но это значит, что скан по времени запуска — твой самый горячий запрос, и он обязан идти по частичному индексу на (run_at) WHERE state = 'pending', а не по sequential scan растущей таблицы. Прикинь объём: при 500 задач/с на enqueue и хранении 7 дней ты несёшь ~300M строк, так что индекс — и план партиционирования или архивации — не опциональны. Enqueue должен быть одним зафиксированным INSERT, чтобы краш между приёмом и сохранением ничего не терял.

    Критерии готовности
    • Задачи переживают рестарт процесса (зафиксированы в Postgres, никогда в памяти), и ты можешь показать состояния жизненного цикла задачи.
    • Скан по времени запуска использует частичный индекс (ты подтвердил через EXPLAIN, что это index scan, а не seq scan), и ты сформулировал допущения по числу строк и хранению.
    Самопроверка

    Покажи EXPLAIN запроса захвата по мере роста таблицы; senior-ревьюер проверяет, что скан по времени остаётся index scan и что хранение не превратит его в seq scan через месяц.

  2. 02At-least-once: захват, visibility timeout, heartbeat

    Сделай доставку устойчивой к смерти воркера на середине задачи. Воркер захватывает наступившую задачу, переводя её в 'claimed' с истекающей арендой (visibility timeout) — параллельные воркеры не должны схватить ту же строку, поэтому захват — это атомарный UPDATE ... WHERE state='pending' ... RETURNING с FOR UPDATE SKIP LOCKED. Если воркер падает до завершения, аренда истекает и задача снова доступна для захвата: это и есть at-least-once, и это единственная честная гарантия, которую долговечная очередь даёт дёшево. Настрой таймаут под p99 длительности задачи — слишком короткий, и медленный-но-живой воркер теряет задачу, которую запустят дважды; слишком длинный, и реальный краш оставляет задачу залипшей на минуты. Длинные задачи продлевают аренду heartbeat'ом (например, обновление каждые 10с при visibility timeout 30с), чтобы reclaim управлялся живостью, а не фиксированной догадкой.

    Критерии готовности
    • Два воркера, гонящиеся за одной наступившей задачей, никогда не запускают её оба (доказано через FOR UPDATE SKIP LOCKED или эквивалентный атомарный захват при конкурентной нагрузке).
    • У воркера, убитого на середине задачи, задача переисполняется после visibility timeout, а длинная задача удерживает аренду через heartbeat, а не теряется.
    Самопроверка

    При at-least-once что будет, если воркер закончил работу, но убит до отметки задачи как done? Senior-ревьюер проверяет, что ответ — «она перезапустится — для этого и нужен следующий этап», а не «потеряна».

  3. 03Ретраи: экспоненциальный backoff, jitter, dead-letter

    Упавшая задача должна повторяться — но наивные ретраи это самонанесённый сбой. Ретраи с фиксированным интервалом от тысяч задач синхронизируются в thundering herd, который снова насыщает только что упавшую зависимость; экспоненциальный backoff (например, 1с, 2с, 4с, 8с, с потолком 5мин) их разносит, а full jitter (sleep = random(0, backoff)) ломает синхронизацию, чтобы ретраи не стреляли все в один тик. Ограничь бюджет попыток: после N попыток (обычно 5–8) задача — poison и должна уйти в dead-letter, а не повторяться вечно и забивать полосу. DLQ — это запрашиваемая таблица, по глубине которой ты алертишь: растущая глубина DLQ — твой самый ранний сигнал, что downstream сломан или payload битый.

    Критерии готовности
    • Ретраи используют экспоненциальный backoff с full jitter и потолком задержки, и ты можешь показать, что две падающие задачи не повторяются в унисон.
    • Задача, упавшая N раз, попадает в dead-letter таблицу (а не повторяется вечно), и глубина DLQ запрашиваема для алертинга.
    Самопроверка

    Покажи расписание backoff и формулу jitter; senior-ревьюер проверяет, что jitter полный (random(0, cap)), а не фиксированный, и что потолок попыток ведёт в DLQ, а не в бесконечный цикл ретраев.

  4. 04Cron и отложенные задачи: расписание, clock skew, двойной запуск

    Добавь повторяющиеся расписания — и ты импортируешь сложнейшую проблему системы: время. Запись cron (например, '*/5 * * * *') должна материализоваться в конкретные строки run_at, а сам материализатор ради доступности работает больше чем на одном узле — поэтому два планировщика могут вычислить одно и то же время запуска и поставить одну минуту дважды. Clock skew между узлами делает «уже пора?» неоднозначным у границы, а переход на летнее время или leap second может пропустить или удвоить окно. Лекарство — сделать запуск детерминированным и дедуплицированным: выведи ключ на каждое срабатывание (schedule_id + забакетированная метка запуска) и ставь в очередь с unique-ограничением, чтобы двойное вычисление схлопнулось в одну строку. Явно реши политику таймзоны и DST — «запустить в 02:30 по местному» — это вопрос, а не момент, в ночь перевода стрелок вперёд.

    Критерии готовности
    • Cron-расписание материализуется в строки run_at, и два экземпляра планировщика, вычислившие одно время запуска, ставят ровно одну задачу (уникальный ключ на срабатывание, а не best-effort).
    • Ты задокументировал политику таймзоны/DST и можешь объяснить, что делает планировщик на границе перехода вперёд или назад.
    Самопроверка

    Запусти два экземпляра планировщика на одном расписании с разведёнными часами; senior-ревьюер проверяет, что ключ дедупа на срабатывание схлопывает двойной запуск в одну строку и что политика DST записана, а не подразумевается.

  5. 05Эффект exactly-once: dedup-хранилище, транзакционные обработчики

    Доставка at-least-once улажена; теперь сделай дублирующую доставку безвредной. Доставку exactly-once дёшево не получить, но эффект exactly-once — можно: каждый обработчик берёт ключ идемпотентности и записывает его в dedup-хранилище в той же транзакции, что и побочный эффект, так что повтор становится no-op. Тонкость — граница атомарности: если работа и запись dedup-ключа в разных коммитах, краш между ними либо применяет дважды (ключ записан после эффекта, но эффект не идемпотентен), либо повторяет вечно (эффект сделан, ключ не записан). Для эффектов вне Postgres (отправка письма, вызов платёжного API) dedup-запись и локальное состояние должны коммититься вместе, а внешний вызов сам должен быть keyed, иначе откатывайся к паттерну outbox: зафиксируй намерение локально, потом доставляй. Докажи: повтор любой задачи дважды даёт ровно один побочный эффект.

    Критерии готовности
    • Каждый обработчик keyed по ключу идемпотентности, записанному в dedup-хранилище в той же транзакции, что и эффект, и повтор задачи дважды даёт ровно один побочный эффект (доказано тестом).
    • Для эффекта вне Postgres ты можешь назвать границу атомарности (транзакционный dedup или outbox) и объяснить, почему наивное «сделать работу, потом отметить done» применяет дважды при краше.
    Самопроверка

    Пройди окно краша между выполнением работы и записью dedup-ключа; senior-ревьюер проверяет, что эффект и ключ коммитятся атомарно (или через outbox), так что ни двойное применение, ни бесконечный повтор невозможны.

  6. 06Масштаб и выживание: приоритет, честность, наблюдение, инцидент poison-pill

    Запусти под нагрузкой и сделай читаемым, потом сломай. Много воркеров, опрашивающих одну таблицу, конкурируют за одни горячие строки, поэтому захваты должны использовать SKIP LOCKED, а ты должен рассудить интервал опроса против латентности (частый опрос жжёт CPU БД; редкий добавляет задержку). Добавь приоритетные полосы — высокоприоритетные задачи прыгают через очередь, — но чистый приоритетный столбец уморяет низкоприоритетную работу голодом, поэтому повышай по возрасту (взвешенная честная очередь или старейшая-подходящая-в-бюджете), чтобы ничто не ждало вечно. Сними RED на воркер-цикле: rate захватов, rate отказов и длительность задачи, плюс глубину очереди и возраст старейшей-незахваченной как SLO, который реально важен. Затем проведи инцидент: poison-pill задача (payload, крашащий обработчик каждую попытку) вместе с переводом часов, выстреливающим thundering herd наступивших задач разом — воркеры крэш-лупят на poison-задаче, пока herd забивает полосу. Обнаружь по возрасту старейшей-незахваченной и глубине DLQ, смягчи вживую (быстро уведи poison в DLQ, честность держит полосу в движении), затем напиши корневую причину.

    Критерии готовности
    • Много воркеров масштабируют пропускную способность захвата без двойного захвата (SKIP LOCKED), а приоритетные полосы повышают по возрасту, так что низкоприоритетные задачи не голодают при устойчивой высокоприоритетной нагрузке.
    • Дашборд показывает rate захвата/отказов, длительность задачи, глубину очереди и возраст старейшей-незахваченной, привязанные к SLO; трейс или структурный лог даёт локализовать медленный обработчик.
    • Ты воспроизвёл инцидент poison-pill + thundering-herd, смягчил его вживую (быстрый увод в DLQ + честность) и написал пост-мортем, чья корневая причина называет механизм (крэш-луп poison / синхронный запуск), а не просто «очередь забилась».
    Самопроверка

    Вставь корневую причину пост-мортема и фикс честности/увода poison; senior-ревьюер проверяет, что SLO — это возраст старейшей-незахваченной (а не просто пропускная способность), что приоритет повышается по возрасту и что poison-pill назван триггером — а не просто «высокая латентность».

Рубрика

Джуниор Миддл Сеньор
Материализация cron и предотвращение двойного запуска Запись cron опрашивается в цикле и задача ставится в очередь при срабатывании; два экземпляра планировщика, работающих параллельно, могут оба сработать на одну минуту и поставить одну задачу дважды. Каждое срабатывание cron получает детерминированный ключ (schedule_id + забакетированная метка запуска) и ставится в очередь с уникальным ограничением, так что два планировщика, вычислившие одно время, схлопываются в одну строку. Политика таймзоны/DST записана — «запустить в 02:30 по местному» — задокументированный вопрос, а не неявный момент, в ночь перевода стрелок вперёд. Ты можешь описать, что делает планировщик, когда DST-разрыв пропускает запланированную минуту (пропустить или запустить на следующей доступной минуте), и какая функция postgres clock используется и почему она стабильна между репликами.
Доставка at-least-once и visibility timeout Воркер захватывает задачу и отмечает её как done после завершения; если воркер падает на середине, задача не повторяется — она навсегда застревает в 'claimed'. Захват — это атомарный UPDATE ... WHERE state='pending' ... RETURNING с FOR UPDATE SKIP LOCKED; visibility timeout погашает аренду и делает задачу захватываемой снова; длинные задачи продлевают аренду через heartbeat, а не теряются. Visibility timeout настроен против p99 длительности задачи — значение короче p99 заставляет медленного-но-живого воркера терять задачу, которую запускают дважды, создавая дубль, который следующий этап должен сделать безвредным. Ты можешь назвать худшее окно повторного выполнения (max_attempts × (visibility_timeout + max_backoff)) и обосновать выбранный таймаут против этого числа.
Backoff, jitter и маршрутизация в dead-letter Упавшие задачи повторяются с фиксированным интервалом; множество падающих задач стреляют в один тик и создают thundering herd против той же сломанной зависимости. Ретраи используют экспоненциальный backoff с full jitter (sleep = random(0, cap)) и потолком попыток; задача, исчерпавшая попытки, уходит в dead-letter таблицу, а не обратно в pending. Ты можешь показать, что две задачи, упавшие в один момент, не повторяются синхронно при данной формуле jitter (гистограмма или diff значений next_run_at в логах). Глубина DLQ запрашиваема и привязана к алерту, так что растущая DLQ — самый ранний сигнал о том, что downstream сломан или класс payload некорректен, а не что-то обнаруживаемое через дни по жалобе пользователя.
Эффект exactly-once и идемпотентные обработчики Обработчики запускают побочный эффект и затем отмечают задачу как done; краш между двумя шагами либо перезапускает эффект, либо оставляет задачу застрявшей — механизма дедупа нет. Каждый обработчик записывает ключ идемпотентности в той же транзакции, что и побочный эффект, так что повторно доставленная задача — это no-op. Для эффектов вне Postgres dedup-запись и локальное состояние коммитятся атомарно (или через outbox), так что ни двойное применение, ни бесконечный повтор невозможны. Chaos-тест — убийство воркеров в случайных точках — даёт ноль потерянных задач и ноль дублированных побочных эффектов на прогоне ≥ 1000 задач смешанных типов. Ты можешь пройти окно краша для не-Postgres-эффекта (отправка письма, платёжный вызов) и назвать точный порядок коммитов, делающий его безопасным: dedup-ключ-в-той-же-транзакции против outbox-публикуй-потом-доставляй, и когда каждый из них правильный выбор.
Эталонный разбор (спойлер)

Почему visibility timeout вместо простой блокировки: распределённая блокировка, удерживаемая упавшим процессом, удерживается вечно — никакой координатор не знает, что он мёртв. Visibility timeout истекает автоматически, делая задачу захватываемой снова после ограниченной задержки. Компромисс — доставка at-least-once: одна задача может выполниться дважды, если воркер завершает работу, но падает до отметки done, именно поэтому обработчики должны быть идемпотентными.

Full jitter не опционален: ретраи с фиксированным интервалом от тысяч упавших задач синхронизируются в thundering herd, который снова насыщает только что упавшую зависимость, превращая транзиентный сбой в затяжной. random(0, min(cap, base × 2^попытка)) гарантирует равномерное распределение времён ретрая в окне backoff, так что нагрузка на восстанавливающуюся зависимость распределяется, а не концентрируется.

Дедупликация cron должна быть детерминированной, а не best-effort: cron, срабатывающий сравнением «now >= next_run_at» в двух экземплярах планировщика, поставит одну минуту дважды всякий раз, когда оба просыпаются в одну секунду. Уникальный ключ на срабатывание (хеш schedule_id + floor(fire_ts, минута)) схлопывает гонку до единственного конфликта INSERT, что дёшево и корректно.

Граница атомарности для эффекта exactly-once: если обработчик записывает побочный эффект, а затем dedup-ключ в двух отдельных коммитах, краш между ними либо применяет дважды (эффект сделан, ключ не записан → повтор снова запускает эффект), либо повторяет вечно (ключ записан, эффект не сделан → повтор — no-op, но эффект так и не случился). Два записи должны быть единым коммитом внутри Postgres, или намерение должно пройти через паттерн outbox для внешних эффектов.

Сделай по-сеньорски

  • Гарантируй эффект exactly-once для эффекта вне Postgres (письмо/платёж) через транзакционный outbox: зафиксируй намерение + dedup-ключ локально, доставляй из outbox с ключом и докажи, что краш в любой точке не применяет дважды и не теряет эффект.
  • Реализуй взвешенную честную очередь с повышением по возрасту, чтобы высокоприоритетные полосы прыгали через очередь, но низкоприоритетные задачи не голодали — проверь нагрузочным тестом, что ожидание старейшей низкоприоритетной задачи остаётся ограниченным при устойчивом высокоприоритетном давлении.
  • Добавь полосу fan-out с ограничением частоты: один триггер разворачивается в тысячи дочерних задач, доставляемых под глобальным потолком token-bucket, чтобы fan-out не устроил thundering herd downstream'у — замерь, что QPS downstream остаётся под его лимитом.
  • Шардируй таблицу очереди по run_at (или хешу) с политикой архивации/удаления партиций, чтобы скан по времени оставался index scan при 300M+ строк, а хранение было дешёвым drop партиции, а не выматывающим DELETE.

Навыки

cron expression parsingvisibility timeout + heartbeatexponential backoff with jitteridempotency keysdead-letter queue

Рекомендуемый стек

nodepostgreshono