sqlpostgresqlconcurrencylocking

Очередь задач в PostgreSQL через FOR UPDATE SKIP LOCKED

Как заставить десяток воркеров безопасно разбирать задачи из таблицы-очереди, не блокируя друг друга — с помощью FOR UPDATE SKIP LOCKED в PostgreSQL.

10 мин чтенияСправочникsql · postgresql · concurrency · locking · job-queue

Иногда для проекта не хочется сразу поднимать отдельный брокер очередей вроде Kafka или RabbitMQ. Нужно что-то проще: отправить письмо после регистрации, пересчитать отчёт, сгенерировать PDF, обработать заказ, вызвать внешний API.

Если задач не миллионы в секунду, а обычная рабочая фоновая очередь для приложения, PostgreSQL может справиться сам. Причём не «на честном слове», а с нормальными блокировками, транзакциями и защитой от гонок.

Главная строчка здесь такая:

FOR UPDATE SKIP LOCKED

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

Разберём всё на понятном примере: у нас есть таблица jobs, где каждая строка — фоновая задача.

Зачем вообще делать очередь в базе

Допустим, пользователь оформил заказ. Сразу после этого нужно:

  • отправить письмо с подтверждением;
  • обновить отчёт;
  • передать заказ во внешнюю систему;
  • создать PDF-чек;
  • отправить уведомление менеджеру.

Не всегда хочется делать это прямо в основном запросе пользователя. Пользователь ждёт страницу, а генерация PDF или внешний API могут тормозить.

Поэтому приложение кладёт задачу в очередь:

Нужно обработать заказ 1001.

А отдельный процесс — воркер — потом забирает эту задачу и выполняет.

Простейшая очередь в PostgreSQL — это обычная таблица.

Таблица задач

Создадим таблицу jobs.

CREATE TABLE jobs (
    id bigserial PRIMARY KEY,
    order_id bigint NOT NULL,
    status text NOT NULL DEFAULT 'pending',
    run_after timestamptz NOT NULL DEFAULT now(),
    attempts int NOT NULL DEFAULT 0,
    locked_by text,
    created_at timestamptz NOT NULL DEFAULT now()
);

Что здесь хранится:

  • id — номер задачи;
  • order_id — заказ, с которым связана задача;
  • status — состояние задачи;
  • run_after — раньше этого времени задачу брать нельзя;
  • attempts — сколько раз уже пытались выполнить задачу;
  • locked_by — какой воркер забрал задачу;
  • created_at — когда задача появилась.

Обычно у задачи бывают такие статусы:

pending
processing
done
failed

Добавим несколько задач:

INSERT INTO jobs (order_id)
VALUES
    (1001),
    (1002),
    (1003);

Теперь у нас есть три задачи в очереди.

Наивный вариант: почему он ломается

На первый взгляд воркер может делать так:

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
ORDER BY id
LIMIT 1;

Он нашёл первую задачу, например id = 1, а потом вторым запросом пометил её как взятую в работу:

UPDATE jobs
SET status = 'processing'
WHERE id = 1;

Для одного воркера это вроде бы работает. Но как только воркеров становится два, появляется гонка.

Представим две сессии.

Первый воркер выполняет:

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
ORDER BY id
LIMIT 1;

Он получает:

id | order_id
---+---------
1  | 1001

Почти одновременно второй воркер выполняет тот же запрос и тоже получает:

id | order_id
---+---------
1  | 1001

Почему так произошло? Потому что первый воркер ещё не успел обновить строку. Между SELECT и UPDATE есть маленькое окно, куда успел пролезть второй воркер.

Итог неприятный: оба воркера начинают обрабатывать один заказ. Письмо может уйти дважды, PDF может сгенерироваться дважды, внешний API может получить два одинаковых запроса.

Это и есть классическая гонка.

Что делает FOR UPDATE

PostgreSQL умеет блокировать строки прямо во время чтения.

Для этого к SELECT добавляют FOR UPDATE.

BEGIN;

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
ORDER BY id
FOR UPDATE
LIMIT 1;

Такой запрос не просто читает строку. Он говорит:

Я выбрал эту строку и собираюсь её менять. До конца моей транзакции другим лучше её не трогать.

Строка остаётся заблокированной до завершения транзакции:

COMMIT;

или:

ROLLBACK;

Это уже защищает нас от ситуации, когда два воркера одновременно спокойно читают одну и ту же задачу.

Но есть новая проблема.

Если первый воркер заблокировал задачу id = 1, второй воркер при попытке взять ту же строку будет ждать. Он не пойдёт к следующей задаче, хотя id = 2 и id = 3 свободны.

Для очереди это плохо. Нам не нужно, чтобы воркеры стояли в пробке. Нам нужно, чтобы каждый взял свободную задачу и пошёл работать.

И вот здесь появляется SKIP LOCKED.

Что делает SKIP LOCKED

SKIP LOCKED означает:

Если строка уже заблокирована другой транзакцией, не жди её. Просто пропусти и ищи следующую подходящую строку.

Именно это нужно для очереди задач.

BEGIN;

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
  AND run_after <= now()
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1;

COMMIT;

Теперь, если два воркера одновременно полезут в очередь:

  • первый заблокирует задачу id = 1;
  • второй не будет ждать id = 1;
  • второй пропустит заблокированную строку и возьмёт id = 2.

Получается простая и красивая схема: много воркеров читают одну таблицу, но каждый получает свою задачу.

Полный вариант с транзакцией

Один только SELECT мало полезен. После выбора задачи нужно сразу пометить её как взятую в работу.

BEGIN;

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
  AND run_after <= now()
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1;

UPDATE jobs
SET status = 'processing',
    locked_by = 'worker-7',
    attempts = attempts + 1
WHERE id = 1;

COMMIT;

Что здесь происходит:

  1. Начинаем транзакцию через BEGIN.
  2. Ищем одну готовую задачу.
  3. Блокируем выбранную строку через FOR UPDATE.
  4. Не ждём чужие заблокированные строки благодаря SKIP LOCKED.
  5. Помечаем задачу как processing.
  6. Завершаем транзакцию через COMMIT.

После COMMIT строка уже не просто заблокирована, а реально обновлена. У неё статус processing, поэтому другие воркеры больше не выберут её по условию:

status = 'pending'

Почему выборку и UPDATE нужно держать вместе

Важно: выбор задачи и перевод в processing должны быть одной атомарной операцией.

Если вы сделали SELECT, потом отпустили транзакцию, а потом отдельным действием сделали UPDATE, гонка может вернуться.

Правильная идея такая:

Забрал строку, сразу пометил её как занятую, зафиксировал транзакцию.

Не нужно держать задачу только «в памяти приложения». Состояние должно быть записано в базе.

Более удобный способ: взять задачу одним запросом

В реальном коде обычно не хотят делать отдельно SELECT, потом руками подставлять id в UPDATE.

Гораздо удобнее сделать всё одним запросом через CTE.

WITH picked AS (
    SELECT id
    FROM jobs
    WHERE status = 'pending'
      AND run_after <= now()
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 1
)
UPDATE jobs AS j
SET status = 'processing',
    locked_by = 'worker-7',
    attempts = j.attempts + 1
FROM picked
WHERE j.id = picked.id
RETURNING j.id, j.order_id;

Этот запрос сразу делает несколько вещей:

  1. В CTE picked выбирает свободную задачу.
  2. Блокирует её.
  3. Обновляет статус на processing.
  4. Возвращает воркеру данные задачи через RETURNING.

Воркер получает результат:

id | order_id
---+---------
1  | 1001

И может начинать работу.

Если свободных задач нет, RETURNING ничего не вернёт. Воркер может немного подождать и попробовать снова.

Как брать задачи пачкой

Иногда воркеру удобно брать не одну задачу, а пачку. Например, сразу 10 писем для отправки.

Для этого меняем только LIMIT.

WITH picked AS (
    SELECT id
    FROM jobs
    WHERE status = 'pending'
      AND run_after <= now()
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 10
)
UPDATE jobs AS j
SET status = 'processing',
    locked_by = 'worker-7',
    attempts = j.attempts + 1
FROM picked
WHERE j.id = picked.id
RETURNING j.id, j.order_id;

Теперь один воркер заберёт до 10 задач.

Если параллельно работают другие воркеры, они заберут другие строки. Заблокированные строки будут пропущены.

Это и есть главное преимущество SKIP LOCKED: воркеры не толкаются у одной двери, а быстро разбирают свободные задачи.

Что делать после успешной обработки

Когда воркер выполнил задачу, нужно пометить её как завершённую.

UPDATE jobs
SET status = 'done',
    locked_by = NULL
WHERE id = 1;

Можно дополнительно добавить поля вроде finished_at, result, error_message, если они нужны вашему проекту.

Например:

ALTER TABLE jobs
ADD COLUMN finished_at timestamptz,
ADD COLUMN last_error text;

Тогда успешное завершение можно записать так:

UPDATE jobs
SET status = 'done',
    locked_by = NULL,
    finished_at = now(),
    last_error = NULL
WHERE id = 1;

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

Что делать, если задача упала

Фоновые задачи иногда падают. Внешний сервис недоступен, письмо не отправилось, PDF не сгенерировался, сеть моргнула.

Обычно задачу не стоит сразу выбрасывать. Её возвращают в очередь с задержкой.

UPDATE jobs
SET status = 'pending',
    locked_by = NULL,
    run_after = now() + interval '30 seconds'
WHERE id = 1
  AND attempts < 5;

Так задача снова станет доступна не сразу, а через 30 секунд.

Если попыток уже слишком много, можно перевести задачу в failed.

UPDATE jobs
SET status = 'failed',
    locked_by = NULL
WHERE id = 1
  AND attempts >= 5;

Идея простая:

  • временная ошибка — повторяем позже;
  • слишком много ошибок — отправляем в failed;
  • дальше человек или отдельный процесс разбирается, что произошло.

Почему не надо держать транзакцию открытой во время работы

Очень важный момент: не держите транзакцию открытой, пока воркер реально выполняет долгую задачу.

Плохой сценарий выглядит так:

BEGIN;

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1;

А потом внутри той же транзакции воркер минуту генерирует PDF, ходит во внешний API или отправляет письмо.

Так делать не стоит.

Транзакция должна быть короткой:

  1. Взяли задачу.
  2. Пометили её как processing.
  3. Сделали COMMIT.
  4. Только потом начали долгую работу.

Почему это важно?

Открытая транзакция держит блокировки и мешает базе спокойно обслуживать другие процессы. Если транзакции висят долго, могут копиться проблемы с очисткой старых версий строк, расти нагрузка и ухудшаться работа VACUUM.

Правильный подход:

WITH picked AS (
    SELECT id
    FROM jobs
    WHERE status = 'pending'
      AND run_after <= now()
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 1
)
UPDATE jobs AS j
SET status = 'processing',
    locked_by = 'worker-7',
    attempts = j.attempts + 1
FROM picked
WHERE j.id = picked.id
RETURNING j.id, j.order_id;

Запрос выполнился, транзакция завершилась. После этого воркер спокойно делает долгую работу уже без открытой транзакции.

Почему SKIP LOCKED без транзакции почти бесполезен

FOR UPDATE блокирует строки только до конца транзакции.

Если вы выполняете запрос в режиме autocommit, каждый отдельный запрос сам по себе является маленькой транзакцией. Запрос закончился — блокировка сразу отпустилась.

Поэтому такой подход опасен:

SELECT id, order_id
FROM jobs
WHERE status = 'pending'
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1;

Если после этого вы отдельно делаете UPDATE, между запросами снова появляется окно для гонки.

Именно поэтому хороший вариант — один запрос с CTE и UPDATE. Он выбирает и помечает задачу в рамках одной операции.

Индекс для очереди

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

Наш основной запрос ищет задачи так:

WHERE status = 'pending'
  AND run_after <= now()
ORDER BY id

Значит, индекс тоже стоит сделать под этот сценарий.

CREATE INDEX idx_jobs_pending
ON jobs (run_after, id)
WHERE status = 'pending';

Это частичный индекс. В него попадут только строки со статусом pending.

Почему это хорошо?

Готовые, выполняющиеся и упавшие задачи не раздувают индекс. Очередь быстрее находит именно те строки, которые можно брать в работу.

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

Почему порядок может быть не строгим

В запросе мы пишем:

ORDER BY id

Это задаёт понятный порядок: сначала более старые задачи.

Но SKIP LOCKED может этот порядок ослабить.

Представьте:

  • задача id = 1 уже заблокирована другим воркером;
  • задача id = 2 свободна;
  • ваш воркер делает выборку.

Он не будет ждать id = 1. Он пропустит её и возьмёт id = 2.

Для очереди фоновых задач это обычно нормально. Нам важнее, чтобы воркеры не простаивали.

Но если вам нужен строгий порядок выполнения, где задача id = 2 не имеет права начаться раньше id = 1, SKIP LOCKED может не подойти. Для строгого порядка нужна другая архитектура и другие правила синхронизации.

Запомните:

SKIP LOCKED даёт параллельность, но не обещает идеальный FIFO.

Как вернуть зависшие задачи

В реальной жизни воркер может умереть посреди работы. Например, он уже перевёл задачу в processing, сделал COMMIT, а потом приложение упало.

Такая задача может навсегда остаться в статусе processing, если не предусмотреть защиту.

Для этого обычно добавляют поле locked_at.

ALTER TABLE jobs
ADD COLUMN locked_at timestamptz;

При захвате задачи записываем время:

WITH picked AS (
    SELECT id
    FROM jobs
    WHERE status = 'pending'
      AND run_after <= now()
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 1
)
UPDATE jobs AS j
SET status = 'processing',
    locked_by = 'worker-7',
    locked_at = now(),
    attempts = j.attempts + 1
FROM picked
WHERE j.id = picked.id
RETURNING j.id, j.order_id;

А отдельным запросом можно возвращать зависшие задачи обратно в очередь:

UPDATE jobs
SET status = 'pending',
    locked_by = NULL,
    locked_at = NULL,
    run_after = now() + interval '1 minute'
WHERE status = 'processing'
  AND locked_at < now() - interval '10 minutes'
  AND attempts < 5;

Если попыток уже много, переводим в failed:

UPDATE jobs
SET status = 'failed',
    locked_by = NULL,
    locked_at = NULL
WHERE status = 'processing'
  AND locked_at < now() - interval '10 minutes'
  AND attempts >= 5;

Так очередь становится устойчивее: даже если воркер упал, задача не потеряется навсегда.

Минимальный рабочий шаблон очереди

Если собрать всё вместе, базовая схема выглядит так.

Сначала таблица:

CREATE TABLE jobs (
    id bigserial PRIMARY KEY,
    order_id bigint NOT NULL,
    status text NOT NULL DEFAULT 'pending',
    run_after timestamptz NOT NULL DEFAULT now(),
    attempts int NOT NULL DEFAULT 0,
    locked_by text,
    locked_at timestamptz,
    created_at timestamptz NOT NULL DEFAULT now(),
    finished_at timestamptz,
    last_error text
);

Индекс:

CREATE INDEX idx_jobs_pending
ON jobs (run_after, id)
WHERE status = 'pending';

Захват задачи:

WITH picked AS (
    SELECT id
    FROM jobs
    WHERE status = 'pending'
      AND run_after <= now()
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 1
)
UPDATE jobs AS j
SET status = 'processing',
    locked_by = 'worker-7',
    locked_at = now(),
    attempts = j.attempts + 1
FROM picked
WHERE j.id = picked.id
RETURNING j.id, j.order_id;

Успешное завершение:

UPDATE jobs
SET status = 'done',
    locked_by = NULL,
    locked_at = NULL,
    finished_at = now(),
    last_error = NULL
WHERE id = 1;

Повтор после ошибки:

UPDATE jobs
SET status = 'pending',
    locked_by = NULL,
    locked_at = NULL,
    run_after = now() + interval '30 seconds',
    last_error = 'temporary error'
WHERE id = 1
  AND attempts < 5;

Окончательный провал:

UPDATE jobs
SET status = 'failed',
    locked_by = NULL,
    locked_at = NULL,
    last_error = 'too many attempts'
WHERE id = 1
  AND attempts >= 5;

Этого уже достаточно для простой и надёжной очереди внутри PostgreSQL.

Когда PostgreSQL-очереди достаточно

Очередь в PostgreSQL хорошо подходит, когда:

  • у вас уже есть PostgreSQL;
  • задач не слишком много;
  • нужна простая и понятная инфраструктура;
  • хочется транзакционно связать бизнес-данные и постановку задачи;
  • важно не тащить отдельный брокер ради нескольких фоновых процессов.

Например, приложение создаёт заказ и в той же транзакции кладёт задачу на отправку письма:

BEGIN;

INSERT INTO orders (user_id, amount)
VALUES (42, 9900)
RETURNING id;

INSERT INTO jobs (order_id)
VALUES (1001);

COMMIT;

Если транзакция откатилась, не будет ни заказа, ни задачи. Это удобно: база сохраняет согласованность.

Когда лучше взять отдельный брокер

PostgreSQL может быть хорошей очередью, но не надо превращать его в Kafka.

Отдельный брокер стоит рассмотреть, если:

  • задач очень много;
  • нужен высокий поток событий;
  • нужно много разных подписчиков;
  • нужна сложная маршрутизация сообщений;
  • нужно хранить и переигрывать поток событий;
  • очередь становится центральной частью архитектуры.

Для простых фоновых задач PostgreSQL часто закрывает потребность отлично. Для большой событийной платформы лучше использовать специализированные инструменты.

MySQL, Oracle и ClickHouse

В MySQL 8.0+ тоже есть похожий механизм:

SELECT id
FROM jobs
WHERE status = 'pending'
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1;

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

В Oracle тоже есть SKIP LOCKED, и его часто используют для похожих сценариев.

А вот ClickHouse для такой задачи не подходит. Это колоночная аналитическая СУБД, а не OLTP-база для частых точечных обновлений и построчных блокировок. В ClickHouse удобно быстро читать и анализировать большие объёмы данных, но делать очередь фоновых задач лучше в PostgreSQL, MySQL или отдельном брокере.

Главные ошибки

Первая ошибка — делать SELECT, а потом отдельный UPDATE без общей транзакции. Так два воркера могут взять одну и ту же задачу.

Вторая ошибка — держать транзакцию открытой во время долгой работы. Транзакция нужна только для быстрого захвата задачи, а не для генерации PDF или похода во внешний API.

Третья ошибка — забыть индекс. Без частичного индекса по готовым задачам очередь может начать тормозить на росте таблицы.

Четвёртая ошибка — ждать строгий порядок. SKIP LOCKED специально пропускает занятые строки, поэтому идеальный FIFO не гарантируется.

Пятая ошибка — не продумать зависшие задачи. Если воркер умер после перевода задачи в processing, нужен механизм возврата или перевода в failed.

Главное

FOR UPDATE SKIP LOCKED — это удобный механизм PostgreSQL для параллельной обработки очереди задач.

FOR UPDATE блокирует выбранные строки до конца транзакции.

FOR UPDATE

SKIP LOCKED говорит не ждать строки, которые уже заблокированы другим воркером.

SKIP LOCKED

Вместе они позволяют нескольким воркерам безопасно разбирать задачи из одной таблицы.

Лучший практический шаблон — выбирать и помечать задачи одним запросом через CTE:

WITH picked AS (
    SELECT id
    FROM jobs
    WHERE status = 'pending'
      AND run_after <= now()
    ORDER BY id
    FOR UPDATE SKIP LOCKED
    LIMIT 10
)
UPDATE jobs AS j
SET status = 'processing',
    locked_by = 'worker-7',
    locked_at = now(),
    attempts = j.attempts + 1
FROM picked
WHERE j.id = picked.id
RETURNING j.id, j.order_id;

Так воркер сразу получает свои задачи, а другие воркеры не берут те же строки.

Для простой очереди писем, отчётов, PDF и фоновой обработки PostgreSQL часто достаточно. Главное — держать транзакции короткими, поставить правильный индекс, обрабатывать ошибки и помнить: SKIP LOCKED даёт удобную параллельность, но не строгий порядок до последней строки.

Закрепи на практике

Решай задачи в SQL-тренажёре с мгновенной проверкой и подсказками.

Открыть тренажёр