Материалы / Истории / Как избежать дублей при параллельной обработке
История

Как избежать дублей при параллельной обработке

Как несколько воркеров могут безопасно разбирать очередь PostgreSQL с FOR UPDATE SKIP LOCKED и когда лучше использовать lease-поля.

В переговорке идет обсуждение новой доработки.

Таня (тимлид):

— Эту новую обработку поручений нужно обязательно распараллелить. Без дублей. Ваши предложения?

— Предлагаю добавить в таблицу payments поле locked_by. — оживилась Катя. — И каждый процесс будет обновлять это поле вот так:

BEGIN;
-- Шаг 1: Читаем не обработанные данные.
SELECT id FROM payments
WHERE status = 'new' AND locked_by IS NULL
-- Шаг 2: Захватываем запись
UPDATE payments
SET locked_by = 'worker_123'
WHERE id = <выбранная запись>;
COMMIT; -- Чтобы другие процессы увидели, что запись обрабатывается!
-- Шаг 3: Основная обработка ...
END;

— Мне не нравится промежуточный COMMIT — нахмурился Сергей — Если процесс упадёт на третьем шаге, запись останется заблокированной. Придется следить за “зависшими” платежками.

— И как быть с конкурентным доступом? — спросила Оля — Несколько процессов могут прочитать одинаковые записи с NULL до того, как один из них ее захватит!”

— Ребята, зачем изобретать велосипед? В Postgres уже есть готовые инструменты для этого. — перехватил инициативу Макс. — Вот так это будет выглядеть:

BEGIN;
-- Берём нужное число незаблокированных записей
SELECT * FROM payments
WHERE status = 'new'
ORDER BY created_at -- если нам важно для FIFO
LIMIT 10
FOR UPDATE SKIP LOCKED; -- Магия здесь!
-- Обрабатываем...
COMMIT; -- Записи автоматически разблокируются при завершении транзакции
END;

— И никаких флагов? — удивилась Катя. — Но как это работает?

— Тут простая магия — ответил Макс:

  1. FOR UPDATE сразу блокирует прочитанные строки.

  2. SKIP LOCKED пропускает уже заблокированные записи (они просто не попадут в результаты следующих запросов, пока действует блокировка)

  3. LIMIT определяет количество блокируемых записей для каждого воркера.

— Чтение записей и блокировка выполняются одной командой и в одной транзакции, этим обеспечивается конкурентный доступ. К тому же нет риска “зависших” записей — если процесс упадёт, транзакция откатится автоматически и записи будут разблокированы. — подвел итог Макс.

— Хорошо! — похвалила Таня. — Есть какие-то особенности или ограничения, которые следует учесть?

— Конечно, они есть всегда — откликнулся Макс и начал загибать пальцы:

  1. SKIP LOCKED намеренно пропускает уже заблокированные строки. Это полезно для очереди, но даёт неполное и несогласованное представление данных, поэтому не подходит для аналитического запроса.

  2. Блокирующие clauses нельзя применять там, где возвращённую строку нельзя однозначно связать с исходной: например, после GROUP BY, DISTINCT, агрегатов или некоторых set operations. В сложном JOIN явно указывайте целевую таблицу: FOR UPDATE OF payments SKIP LOCKED.

  3. Если воркеры блокируют несколько таблиц в разном порядке, возможны deadlock. Порядок захвата должен быть единым.

  4. Не держите row lock часами во время внешнего API-вызова. Для долгой обработки лучше коротко присвоить задаче lease (locked_by, locked_until), зафиксировать транзакцию и сделать обработчик идемпотентным. Очередь также нуждается в стратегии против starvation старых задач.

— Подожди, Макс. — задумалась Оля. — А что если обработка платежа занимает не 5 секунд, а 5 часов? Или если нужен ручной контроль менеджера?

— Да, для долгих операций SKIP LOCKED - не лучший выбор. — кивнул Макс. — Тогда лучше через поле locked_by. Для надежности тогда добавить поле locked_until, которое обновлять Heartbeat-механизмом. Это позволит отслеживать зависшие задачи.

— Не отвлекайтесь. — остановила их Таня. — у нас обработка каждого платежа будет очень быстрой, так что решено: используем FOR UPDATE SKIP LOCKED, быстро проверяем на тестовом стенде и идем в продакшн. Главное — ничего не сломать. Макс, проверка на тебе. Переходим к следующему пункту.