Как избежать дублей при параллельной обработке
Как несколько воркеров могут безопасно разбирать очередь PostgreSQL с FOR UPDATE SKIP LOCKED и когда лучше использовать lease-поля.
PostgreSQLБлокировкиКонкурентностьФоновые задачиВ переговорке идет обсуждение новой доработки.
Таня (тимлид):
— Эту новую обработку поручений нужно обязательно распараллелить. Без дублей. Ваши предложения?
— Предлагаю добавить в таблицу 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;— И никаких флагов? — удивилась Катя. — Но как это работает?
— Тут простая магия — ответил Макс:
-
FOR UPDATEсразу блокирует прочитанные строки. -
SKIP LOCKEDпропускает уже заблокированные записи (они просто не попадут в результаты следующих запросов, пока действует блокировка) -
LIMITопределяет количество блокируемых записей для каждого воркера.
— Чтение записей и блокировка выполняются одной командой и в одной транзакции, этим обеспечивается конкурентный доступ. К тому же нет риска “зависших” записей — если процесс упадёт, транзакция откатится автоматически и записи будут разблокированы. — подвел итог Макс.
— Хорошо! — похвалила Таня. — Есть какие-то особенности или ограничения, которые следует учесть?
— Конечно, они есть всегда — откликнулся Макс и начал загибать пальцы:
-
SKIP LOCKEDнамеренно пропускает уже заблокированные строки. Это полезно для очереди, но даёт неполное и несогласованное представление данных, поэтому не подходит для аналитического запроса. -
Блокирующие clauses нельзя применять там, где возвращённую строку нельзя однозначно связать с исходной: например, после
GROUP BY,DISTINCT, агрегатов или некоторых set operations. В сложномJOINявно указывайте целевую таблицу:FOR UPDATE OF payments SKIP LOCKED. -
Если воркеры блокируют несколько таблиц в разном порядке, возможны deadlock. Порядок захвата должен быть единым.
-
Не держите row lock часами во время внешнего API-вызова. Для долгой обработки лучше коротко присвоить задаче lease (
locked_by,locked_until), зафиксировать транзакцию и сделать обработчик идемпотентным. Очередь также нуждается в стратегии против starvation старых задач.
— Подожди, Макс. — задумалась Оля. — А что если обработка платежа занимает не 5 секунд, а 5 часов? Или если нужен ручной контроль менеджера?
— Да, для долгих операций SKIP LOCKED - не лучший выбор. — кивнул Макс. — Тогда лучше через поле locked_by. Для надежности тогда добавить поле locked_until, которое обновлять Heartbeat-механизмом. Это позволит отслеживать зависшие задачи.
— Не отвлекайтесь. — остановила их Таня. — у нас обработка каждого платежа будет очень быстрой, так что решено: используем FOR UPDATE SKIP LOCKED, быстро проверяем на тестовом стенде и идем в продакшн. Главное — ничего не сломать. Макс, проверка на тебе. Переходим к следующему пункту.