infra · advanced · 7d
Очередь задач at-least-once
Собери долговечную очередь задач на Postgres с visibility timeout и идемпотентными консьюмерами, чтобы упавший воркер не терял задачу.
Результат
Очередь на Postgres, где каждая поставленная задача выполняется хотя бы раз при крашах воркеров (SKIP LOCKED + visibility timeout + heartbeat), дубликаты — no-op через атомарные ключи идемпотентности, ядовитые задачи уходят в DLQ после N попыток, а дашборд показывает throughput захвата, перепостановки visibility и глубину DLQ.
Этапы
0/5 · 0%- 01Захват задач без двойного забора
Сделай захват атомарным, чтобы два воркера никогда не брали одну задачу. Наивный паттерн — SELECT ожидающей строки, затем UPDATE в in-flight двумя операторами — это гонка: между SELECT и UPDATE другой воркер может SELECTнуть ту же строку, и оба её UPDATEнут, дважды захватив одну задачу. Почини, схлопнув выборку и блокировку в один оператор: `SELECT ... FOR UPDATE SKIP LOCKED` внутри транзакции, которая также UPDATEит строку в in-flight и RETURNит её. FOR UPDATE блокирует строку, SKIP LOCKED заставляет второго воркера пропустить заблокированную строку вместо блокировки (простой FOR UPDATE сериализовал бы всех воркеров на горячей таблице очереди, добавляя латентность ожидания блокировки пропорционально числу воркеров). Докажи: запусти двух конкурентных воркеров, опрашивающих одну очередь со 100 ожидающими задачами, и убедись, что они никогда не получают один id задачи и что пропущенная строка немедленно доступна другому поллеру — throughput захвата масштабируется с числом воркеров, а не с ожиданиями блокировок.
Критерии готовности- Два параллельных воркера, опрашивающих одну очередь со 100 ожидающими задачами, никогда не получают одну строку; пропущенная по SKIP LOCKED строка подбирается другим воркером за один интервал опроса — проверено конкурентным тестом захвата.
- Захват — одна транзакция: SELECT FOR UPDATE SKIP LOCKED + UPDATE в in-flight + RETURNING, без отдельного окна гонки SELECT-затем-UPDATE.
Опирается наСамопроверка
Покажи SQL захвата и конкурентный тест двух воркеров с id задач. Senior-ревьюер проверяет, что захват — одна транзакция с SKIP LOCKED (не простым FOR UPDATE), и спрашивает, почему простой FOR UPDATE сериализовал бы воркеров.
- 02Перепостановка задач от умерших воркеров
Добавь visibility timeout, чтобы задача, чей воркер умер до ack, не терялась — становилась снова доступной для захвата по истечении таймаута вместо вечного зависания в 'claimed'. Таймаут — самая чувствительная ручка настройки очереди, и сеньорная ошибка — ставить его интуитивно: слишком короткий — медленный, но живой воркер гонится за собственной задачей (задача перепоставлена и обрабатывается дважды параллельно, риск дублирующихся эффектов); слишком длинный — реальный краш стопорит полосу на минуты, прежде чем другой воркер сможет её подобрать. Правильное значение чуть выше p99 длительности задачи, не p50 и не консервативный 10x. Реализуй: захваченные задачи несут `visible_at = now + visibility_timeout`; sweep-процесс или предикат захвата `WHERE visible_at < now()` перепоставляет просроченные задачи. Докажи краш-тестом: убей воркера на середине задачи (SIGKILL до ack) и покажи, что задача становится доступной ровно по истечении visibility timeout, не потеряна и не немедленно.
Критерии готовности- Убийство воркера на середине задачи до ack делает задачу снова доступной ровно по истечении visibility timeout (не потеряна, не немедленно) — проверено краш-тестом до ack с проверкой тайминга перепостановки.
- Значение таймаута задокументировано с p99 длительности задачи, два режима отказа (слишком короткий — гонка с живыми воркерами, слишком длинный — стопор полосы) указаны с обоснованием выбранного значения.
Опирается наСамопроверка
Убей воркера до ack и покажи, что задача появляется после таймаута. Senior-ревьюер проверяет, что таймаут выставлен по p99 (не p50/10x), и просит назвать оба режима отказа, если бы он был в 2 раза короче или в 5 раз длиннее.
- 03Продление аренды heartbeat для длинных задач
Позволь длинным задачам переживать visibility timeout, не будучи украденными. Задача, чья реальная длительность превышает таймаут, иначе была бы перепоставлена, пока ещё выполняется — два воркера тогда обрабатывали бы одну задачу параллельно, превращая at-least-once в at-least-twice с перекрывающимися эффектами. Фикс — heartbeat, продлевающий аренду: воркер периодически UPDATEит `visible_at = now + visibility_timeout` на подинтервале таймаута (например, каждую timeout/3), пока задача in-flight. Продление должно быть атомарным и не продлевать задачу, уже перепоставленную (проверь `WHERE status='claimed' AND visible_at` всё ещё совпадает). Докажи: запусти задачу, спящую 3× visibility timeout с heartbeat, и покажи, что она никогда не перепоставлена; убей heartbeat на середине и покажи перепостановку через один таймаут. Задокументируй, почему увеличение самого таймаута (вместо heartbeat) неверно — оно бы стопорило реальные краши на всю длинную длительность.
Критерии готовности- Задача, спящая 3× visibility timeout с heartbeat каждую timeout/3, никогда не перепоставлена; убийство heartbeat на середине делает её доступной через один таймаут — оба проверены таймированными тестами.
- Продление аренды атомарно (проверяет status/visible_at), выбор интервала и почему 'просто увеличь таймаут' неверно — задокументированы.
Опирается наСамопроверка
Покажи задачу 3×-таймаут с heartbeat, никогда не украденную, и одну где heartbeat остановлен и она перепоставлена. Senior-ревьюер проверяет условность продления по status/visible_at и спрашивает, почему увеличение самого таймаута стопорило бы реальные краши.
- 04Сделай повторную доставку no-op
Сделай консьюмер идемпотентным, чтобы повторно доставленная задача — будь то из visibility timeout или потери heartbeat — давала ровно один эффект. Единственный честный способ — записать ключ идемпотентности в той же транзакции, что и бизнес-эффект: `INSERT dedup_key` + `UPDATE business_table` в одном коммите. Если они в разных коммитах, краш между ними либо применяет дважды (эффект закоммичен, ключ дедупа нет — следующая доставка применяет снова), либо повторяет вечно (ключ дедупа закоммичен, эффекта нет — следующая доставка ошибочно считается выполненной). Натуральный ключ идемпотентности задачи (например, `job_id` или `order_id`) хранится в таблице дедупа с unique-ограничением; повторная доставка с тем же ключом попадает в ограничение и становится no-op. Докажи: обработай одну задачу дважды с тем же ключом и убедись в ровно одном бизнес-эффекте; сымитируй краш между эффектом и вставкой дедупа (провалив транзакцию) и покажи, что ничто не видимо.
Критерии готовности- Обработка одной задачи дважды (тот же ключ идемпотентности) даёт ровно один бизнес-эффект; второй проход — no-op через unique-ограничение дедупа — проверено тестом двойной доставки.
- Ключ дедупа и бизнес-эффект в одной транзакции; краш между ними (симулированный откатом) не оставляет ни того ни другого видимым — граница атомарности протестирована и задокументирована.
Опирается наСамопроверка
Покажи двойную доставку с одним эффектом и откат между вставкой дедупа и бизнес-записью, не оставляющий ничего. Senior-ревьюер проверяет, что оба в одной транзакции, и спрашивает, что сломается при двух коммитах.
- 05Dead-letter, backoff и хаос
Добавь прод-сетку безопасности: ядовитая задача, всегда падающая, не должна крутиться вечно и занимать голову очереди, блокируя каждую последующую задачу. После N неудач (например, 5) с экспоненциальным backoff и jitter (`delay = base * 2^attempt + rand(0, jitter)`) задача уходит в dead-letter таблицу вместо перепостановки — рост глубины DLQ — самый ранний сигнал, что что-то выше по цепи сломано или payload битый, и алертить надо на него, а не только на латентность. Сделай очередь наблюдаемой: снимай метрики throughput захвата, rate перепостановок visibility, глубину DLQ и p50/p99 длительности задач. Затем прогони хаос-тест: убивай воркеров в каждой точке краша (до захвата, во время задачи, до ack, во время вставки дедупа) под конкурентной нагрузкой и убедись в ноле потерянных задач и ноле небезопасных дублей. Хаос-тест — доказательство, что границы атомарности из этапов 1 и 4 действительно держатся, когда процесс умирает в худший момент.
Критерии готовности- Ядовитая задача, всегда падающая, уходит в DLQ после N попыток с экспоненциальным backoff+jitter (задержки проверены); глубина DLQ — метрика, порог алерта задокументирован.
- Хаос-тест, убивающий воркеров в каждой точке краша под конкурентной нагрузкой, показывает ноль потерянных задач и ноль небезопасных дублей; throughput захвата и глубина DLQ видны на дашборде.
Опирается наСамопроверка
Покажи ядовитую задачу в DLQ после N попыток с задержками backoff и хаос-тест, убивающий воркеров в 4 точках краша с нулём потерянных/небезопасных. Senior-ревьюер проверяет экспоненциальный+jitter backoff (не фиксированный), что DLQ-глубина — сигнал алерта, а хаос покрывает границу атомарности дедупа.
Стартер
fallowlone/skein-projects
projects/at-least-once-queue
- README.md
- src/queue.ts
- test/queue.test.ts
npx degit fallowlone/skein-projects/projects/at-least-once-queue at-least-once-queue Реализуй заглушки, затем гоняй тесты, пока не позеленеют: bun test
Форкни репозиторий и запушь свою работу — workflow grade прогонит тесты и статические проверки на твоих раннерах.
Рубрика
| Джуниор | Миддл | Сеньор | |
|---|---|---|---|
| Атомарность захвата | Воркер выбирает задачу и обновляет её двумя отдельными запросами; при конкурентной нагрузке два воркера изредка захватывают одну строку. | Захват — единственный UPDATE ... WHERE state='pending' RETURNING с FOR UPDATE SKIP LOCKED: конкурирующие воркеры никогда не дублируют захват, а пропущенная строка немедленно доступна следующему поллеру. | Ты можешь объяснить модель конкуренции: SKIP LOCKED масштабируется на множество воркеров без ожидания блокировок, но концентрирует всю ожидающую работу на старейших строках; ты измеряешь потолок пропускной способности захвата и знаешь, когда партиционировать таблицу очереди. |
| Visibility timeout и повторная доставка | Задача упавшего воркера зависает в состоянии 'claimed' до ручного вмешательства; автоматической перепостановки нет. | Sweep-процесс или проверка аренды возвращает в очередь задачи с истёкшим visibility timeout; живой воркер продлевает аренду через heartbeat, чтобы длинные задачи не угонялись. | Ты устанавливаешь таймаут по p99 длительности задачи и можешь сформулировать два режима отказа: слишком короткий — медленный, но живой воркер гонится за собственной задачей; слишком длинный — реальный краш стопорит полосу на минуты — ты документируешь выбранное значение и обоснование. |
| Идемпотентный консьюмер и dead-letter | Повторно доставленные задачи обрабатываются снова, изредка порождая дублирующиеся эффекты; ограничения попыток нет. | Ключ дедупликации, записанный в той же транзакции, что и эффект, делает повторную доставку no-op; после N неудач задача уходит в dead-letter вместо цикличного повтора. | Ты рассуждаешь о границе атомарности: если эффект и ключ дедупа в разных коммитах, краш между ними либо применяет дважды, либо повторяет вечно — твой дизайн делает оба варианта невозможными, и ты доказываешь это хаос-тестом, убивающим воркеров в каждой точке краша. |
Эталонный разбор (спойлер)
Почему at-least-once — это честная базовая гарантия: доставка exactly-once требует распределённой координации, которая либо очень дорога (двухфазный коммит), либо невозможна между разнородными системами. Любая долговечная очередь на едином хранилище может обещать только at-least-once: задача перезапустится, если воркер упадёт до ack, а корректность перекладывается на консьюмер через идемпотентность.
FOR UPDATE SKIP LOCKED как примитив захвата: он объединяет выборку и блокировку в одном операторе, и ни один второй воркер не может увидеть ту же строку; SKIP предотвращает очередь ожидания блокировок — блокирующий FOR UPDATE сериализовал бы всех воркеров на горячей таблице вместо разветвления.
Ловушка настройки visibility timeout: правильный таймаут — чуть выше p99 длительности задачи, не p50 и не консервативный 10x. Слишком короткий — гонка с живыми воркерами; слишком длинный — задачи упавших воркеров стопорят полосу. Heartbeat с обновлением аренды на интервале меньше таймаута — правильный фикс для длинных задач, а не увеличение таймаута.
Глубина dead-letter как самый ранний сигнал: ядовитая задача кладёт каждого воркера, который её касается, и без DLQ она вечно занимает голову очереди, блокируя все последующие задачи. Рост глубины DLQ — первый наблюдаемый симптом того, что что-то выше по цепи сломано или payload битый — алерть на него, а не только на латентность задач.
Сделай по-сеньорски
- Партиционируй таблицу очереди по хешу типа задачи, чтобы throughput захвата масштабировался за потолок single-table SKIP LOCKED; измерь потолок до и после.
- Добавь fair queueing, чтобы всплеск одного тенанта не голодал задачи другого — реализуй per-tenant квоты захвата и докажи изоляцию под нагрузкой.
- Замени поллинг на LISTEN/NOTIFY, чтобы латентность захвата упала с интервала опроса до латентности уведомления; измерь улучшение p99 задержки захвата.