Иногда для проекта не хочется сразу поднимать отдельный брокер очередей вроде 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;
Что здесь происходит:
- Начинаем транзакцию через
BEGIN.
- Ищем одну готовую задачу.
- Блокируем выбранную строку через
FOR UPDATE.
- Не ждём чужие заблокированные строки благодаря
SKIP LOCKED.
- Помечаем задачу как
processing.
- Завершаем транзакцию через
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;
Этот запрос сразу делает несколько вещей:
- В CTE
picked выбирает свободную задачу.
- Блокирует её.
- Обновляет статус на
processing.
- Возвращает воркеру данные задачи через
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 или отправляет письмо.
Так делать не стоит.
Транзакция должна быть короткой:
- Взяли задачу.
- Пометили её как
processing.
- Сделали
COMMIT.
- Только потом начали долгую работу.
Почему это важно?
Открытая транзакция держит блокировки и мешает базе спокойно обслуживать другие процессы. Если транзакции висят долго, могут копиться проблемы с очисткой старых версий строк, расти нагрузка и ухудшаться работа 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 даёт удобную параллельность, но не строгий порядок до последней строки.
Иногда для проекта не хочется сразу поднимать отдельный брокер очередей вроде Kafka или RabbitMQ. Нужно что-то проще: отправить письмо после регистрации, пересчитать отчёт, сгенерировать PDF, обработать заказ, вызвать внешний API.
Если задач не миллионы в секунду, а обычная рабочая фоновая очередь для приложения, PostgreSQL может справиться сам. Причём не «на честном слове», а с нормальными блокировками, транзакциями и защитой от гонок.
Главная строчка здесь такая:
FOR UPDATE SKIP LOCKEDОна позволяет нескольким воркерам одновременно брать задачи из одной таблицы так, чтобы они не хватали одну и ту же строку и не стояли друг за другом в ожидании блокировок.
Разберём всё на понятном примере: у нас есть таблица
jobs, где каждая строка — фоновая задача.Зачем вообще делать очередь в базе
Допустим, пользователь оформил заказ. Сразу после этого нужно:
Не всегда хочется делать это прямо в основном запросе пользователя. Пользователь ждёт страницу, а генерация PDF или внешний API могут тормозить.
Поэтому приложение кладёт задачу в очередь:
А отдельный процесс — воркер — потом забирает эту задачу и выполняет.
Простейшая очередь в 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— когда задача появилась.Обычно у задачи бывают такие статусы:
Добавим несколько задач:
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;Что здесь происходит:
BEGIN.FOR UPDATE.SKIP LOCKED.processing.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;Этот запрос сразу делает несколько вещей:
pickedвыбирает свободную задачу.processing.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 или отправляет письмо.
Так делать не стоит.
Транзакция должна быть короткой:
processing.COMMIT.Почему это важно?
Открытая транзакция держит блокировки и мешает базе спокойно обслуживать другие процессы. Если транзакции висят долго, могут копиться проблемы с очисткой старых версий строк, расти нагрузка и ухудшаться работа
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может не подойти. Для строгого порядка нужна другая архитектура и другие правила синхронизации.Запомните:
Как вернуть зависшие задачи
В реальной жизни воркер может умереть посреди работы. Например, он уже перевёл задачу в
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 хорошо подходит, когда:
Например, приложение создаёт заказ и в той же транзакции кладёт задачу на отправку письма:
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 UPDATESKIP 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даёт удобную параллельность, но не строгий порядок до последней строки.