nullinside1 час назадОбъяснить с
Распределённые блокировки в Multi-DC: как перестать дублировать CronJobs
Уровень сложностиСреднийВремя на прочтение9 минОхват и читатели2.2KБлог компании Точка БанкPython*Веб-разработка*Качество кода*Программирование*ТуториалПредставьте: вы настроили CronJob, протестировали локально и выкатили в продакшн. Всё работает как часы. А потом решили добавить отказоустойчивости и развернули приложение во втором дата‑центре. И внезапно ваша задача начала выполняться дважды. Это не баг Kubernetes, а фундаментальная проблема синхронизации.
Меня зовут Тимофей Салтымаков, я бэкенд‑разработчик в Точка Банк. В статье расскажу, как мы решили проблему дублирующихся CronJob'ов в распределённой системе. Прошли путь от advisory locks до собственной библиотеки, и в итоге поняли, что нужно перестать синхронизировать процессы и начать синхронизировать... время. Поехали!
С чего всё началось
Представим, что у нас есть простая Cron‑задача, которая должна запускаться раз в минуту. Например:
• пересчитать агрегаты;
• отправить уведомления;
• получить данные со сторонней площадки.Пока у нас один сервер или один кластер Kubernetes, всё работает хорошо. Мы настраиваем CronJob, выкатываем его в продакшен и задача запускается чётко по расписанию.
Но если приложение развёрнуто сразу в двух дата-центрах, в каждом из которых работает свой кластер Kubernetes, то наша задача выполнится дважды (или трижды, если у вас три независимых кластера). Данные будут дублироваться и лимиты запросов сгорят в два раза быстрее.
При этом с точки зрения Kubernetes всё работает правильно. Каждый кластер существует сам по себе и ничего не знает о соседних. Он видит только собственное расписание и честно запускает задачу в назначенное время. То есть проблема не на уровне Kubernetes, а на уровне требований к задаче.
Немного про крон в Kubernetes
Для начала сделаем небольшое лирическое отступление. Если вы привыкли к классическому Linux, то слово «крон» для вас, скорее всего, означает строку в crontab. Есть один сервер, на нём работает демон cron, который раз в минуту проверяет расписание и запускает нужные задачи.
В Kubernetes используется другая модель. Расписание хранится в объекте CronJob. Когда наступает время запуска, Kubernetes создаёт объект Job, а тот, в свою очередь, поднимает отдельный Pod. После выполнения задачи Pod завершается и удаляется.
В чём ключевое отличие: в обычном приложении несколько задач могут жить внутри одного процесса и пользоваться общими примитивами синхронизации. У CronJob такой возможности нет. Каждый запуск — это отдельный Pod со своей памятью и жизненным циклом. Он ничего не знает о других экземплярах задачи и никак не может с ними договориться напрямую.
Для одного кластера это не создаёт проблем. Но если одинаковый CronJob развёрнут сразу в нескольких независимых кластерах, каждый из них будет создавать собственный Pod по одному и тому же расписанию. Договариваться между собой процессы могут только через внешний общий компонент — PostgreSQL, Redis или другое хранилище.
Первая мысль — локи
Когда одна и та же задача может выполняться одновременно в нескольких местах, обычно используют локи — механизм взаимного исключения. У него есть три типа операций:
• взять— попытаться захватить лок (если получилось — можно выполнить работу, если нет — значит, лок уже занят);
• удерживать— пока процесс владеет локом, остальные не могут его получить;
• отпустить— освободить лок, чтобы следующий процесс мог продолжить работу.SELECT pg_advisory_lock(67); — будет ждать до освобождения лока, чтобы его занять
SELECT pg_try_advisory_lock(67); — true: лок твой / false: занято
SELECT pg_advisory_unlock(67); — отпуститьДля нашей задачи такая схема выглядит подходящей. Если два CronJob стартуют одновременно в разных дата-центрах, один из них должен получить лок, а второй — увидеть, что работа уже выполняется, и завершиться.
В PostgreSQL для этого есть готовый механизм — advisory lock или «консультативный» лок. В отличие от обычных блокировок строк и таблиц, он не связан с какими‑либо данными. Это просто именованный флаг внутри PostgreSQL. База знает о нём только одно: лок сейчас свободен или занят.
Получить такой лок можно буквально одной командой:
SELECT pg_try_advisory_lock(67)Если функция возвращает true — значит, лок успешно захвачен и можно выполнить задачу. Если false — лок уже удерживает другой процесс.
В итоге алгоритм получился предельно простым:
• CronJob стартует.
• Пытается получить advisory lock.
• Если лок удалось захватить — выполняет задачу.
• Если лок уже занят — сразу завершает работу.На бумаге это выглядело почти идеально. К сожалению, довольно быстро выяснилось, что проблема вообще не в самом механизме блокировки.
Где всё сломалось
Самый важный нюанс advisory lock: он живёт ровно столько, сколько живёт соединение с PostgreSQL.
Пока соединение открыто — лок удерживается. Соединение закрылось — и PostgreSQL автоматически освобождает все advisory lock, которые были за ним закреплены.
Для большинства приложений это отличное решение. Если процесс аварийно завершился, никаких «вечных» блокировок не останется — база сама всё почистит.
Но в нашем случае перед PostgreSQL стоял PgBouncer, который работал в режиме transaction pooling. Это означает, что соединения постоянно переиспользовались между разными клиентами. Чтобы удерживать advisory lock на протяжении всей работы CronJob, пришлось бы удерживать и само соединение. В итоге мы упираемся в длинноживущие транзакции, а это плохая идея для массового инфраструктурного решения.
Попытка номер два — LeaseLock
Мы решили не искать другой готовый механизм, а реализовать собственный лок поверх PostgreSQL. Его состояние должно было храниться локально независимо от соединения с базой.
Для этого достаточно обычной таблицы в PostgreSQL:
CREATE UNLOGGED TABLE key_value (
key text PRIMARY KEY,
value text NOT NULL,
expire_at timestamptz NOT NULL
);Каждая строка в таблице представляет один лок:
• key— уникальный идентификатор лока. PRIMARY KEY гарантирует, что второй процесс не вставит такую же строку. На этом свойстве держится вся блокировка;
• value— владелец лока (в формате hostname:pid);
• expire_at— время жизни лока. Пока expire_at не наступил, лок считается действительным. Время мы берём не у Pod'а, а у базы. Это сделано специально, потому что часы на разных серверах почти никогда не идут синхронно. Если каждый Pod будет считать, протух ли лок, по своим локальным часам, один процесс будет считать лок ещё живым, а второй — уже мёртвым. Единый источник времени решает эту проблему.
Что ещё важно учесть
На первый взгляд LeaseLock выглядит довольно просто. Но чтобы механизм корректно работал при конкурентном доступе и мог пережить падение процессов, важно учесть пару нюансов:
• Проверка и захват должны быть атомарными. Нельзя сначала делать SELECT, а потом INSERT — между ними всегда есть небольшой промежуток времени. Если в этот момент второй процесс выполнит те же действия, оба решат, что лок свободен и начнут выполнять задачу. Поэтому мы используем INSERT... ON CONFLICT с проверкой expire_at.
Если записи ещё нет — она создаётся. Если запись существует, но срок её действия уже истёк — она обновляется. Во всех остальных случаях процесс понимает, что лок занят, и завершает работу.
В итоге вся гонка между дата-центрами сводится к одной SQL‑команде, за корректность которой отвечает сама база. Никаких дополнительных блокировок или сложной логики синхронизации не требуется.
def try_acquire_run_lock(
conn: Any,
table_parts: list[str],
key; str,
ttl_seconds: float,
value: str,
) -> tumple[bool, datetime | None]:
t = make_qualified_table_name(table_parts)
query = sql.SQL(
"""
INSERT INTO {t} (key, value, expire_at)
VALUES (%s, %s, NOW() + (%s * INTERVAL '1 second'))
ON CONFLICT (key) DO UPDATE
SET value = EXCLUDED.value,
expire_at = EXCLUDED.expire_at
WHERE {t}.expire_at <= NOW()
RETURNING key, expire_at;
"""
).format(t=t)
with conn.cursor() as cur:
cur.execute(query, (key, value, ttl_seconds))
row = cur.fetchone()
if row:
return True
return False• Используем отдельный фоновый поток — heartbeat.Предположим, что мы установили expire_at на 60 секунд вперёд. Тогда на 61-й секунде лок автоматически станет свободным, хотя задача может всё ещё выполняться. Чтобы этого не произошло, после успешного захвата запускается фоновый поток — heartbeat. Его единственная задача — периодически продлевать expire_at, пока основная задача жива.
Если процесс завершается штатно — он освобождает лок сам. Если контейнер неожиданно падает — heartbeat прекращается, expire_at перестаёт обновляться, и после истечения TTL лок автоматически становится доступен другому процессу.
def renew_lease_query(table_parts% list[str]) -> Any:
t = make_qualified_table_name(table_parts)
return sql.SQL(
"""
UPDATE {t}
SET expire_at = NOW() + (%s * INTERVAL '1 second')
WHERE key = %s
AND value = %s
RETURNING expire_at;
"""
).format(t=t)В результате получился механизм, который может пережить падение процессов, не зависит от времени жизни соединения с PostgreSQL и хорошо работает в нашей инфраструктуре.
Казалось, что задача решена, но нет. Спустя некоторое время мы поняли, что всё это время боролись не с той проблемой.
Как короткая джоба всё сломала
Представим, что у нас есть задача, которая запускается раз в минуту. Сама работа занимает буквально пару секунд — например, нужно поставить сообщение в очередь Celery или отправить запрос во внешний API.
Теперь посмотрим, что происходит в двух дата-центрах:
• Kubernetes почти одновременно запускает CronJob в двух кластерах.
• Pod в первом дата-центре стартует немного раньше.
• Он берёт lock и выполняет работу.
• lock освобождается.
• Через несколько секунд стартует Pod во втором дата-центре.
• Он видит свободный lock, берёт его и выполняет ту же работу заново.Ни на одном этапе не произошло ошибки. Оба процесса работали корректно и не нарушили правила. Но задача всё равно выполнилась дважды.
Дело в том, что Lease lock отвечает на вопрос: «Кто выполняет работу прямо сейчас?» А проблема CronJob связана не с одновременным выполнением, а с самим фактом выполнения. Поэтому в действительности нас интересует другой вопрос: «Выполнялась ли уже задача для этого запуска расписания?» И здесь нужен другой примитив — schedule lock.
ВопросТип локаРеализация у насКто работает сейчас?lease lockLeaseLockВыполнялось ли за это окно?schedule lockScheduleLock
Наше решение — Schedule Lock
Финальным решением нашей задачи стал Schedule Lock. Его идея заключается в том, что синхронизируются не процессы, а запуски окна выполнения.
У любой Cron‑задачи есть расписание. Например, она может выполняться раз в десять минут:
• 12:00–12:10;
• 12:10–12:20;
• 12:20–12:30;
• 12:30–12:40 и так далее.Для каждого такого окна можно вычислить уникальный идентификатор. Если процесс первым создал ключ — значит, именно он выполняет работу.
Все остальные процессы, независимо от того, стартовали они одновременно или спустя несколько секунд, будут пытаться создать тот же самый ключ. И, если он существует, значит окно расписания уже обработано.
Как только наступает следующий запуск CronJob, начинается новое окно. Старый ключ автоматически становится неактуальным, и задача снова может выполниться — один раз для нового интервала.
from croniter import croniter
# now — это NOW() из PostgreSQL, единый для всех дата-центров
next_run_time = croniter(cron_schedule, db_now).get_next(datetime)from dislock import ScheduleLock, LockNotAcquiredError
try:
with ScheduleLock(context):
run_my_job()
except LockNotAcquiredError:
pass # окно уже отработано в другом DC — выходимТеперь вернёмся к нашему сценарию. Вот что получается:
• Первый кластер стартует, вычисляет текущее окно, создаёт ключ и начинает работу.
• Через несколько секунд стартует второй кластер, вычисляет то же самое окно, пытается занять тот же ключ, получает отказ и завершает работу.
• Даже если работа в первом кластере уже завершилась, повторного выполнения не происходит, так как это окно считается обработанным.Таким образом, мы перестали синхронизировать процессы и начали синхронизировать время. Для большинства Cron‑задач этого достаточно. Исключение составляют сложные сценарии, когда время выполнения задачи больше, чем окно выполнения. Например, если задача запускается каждые пять минут, но сама при этом может выполняться десять минут.
Тут одного ScheduleLock будет недостаточно, нужно использовать связку ScheduleLock + LeaseLock. Первый гарантирует, что каждое окно расписания будет обработано только один раз, а второй — что в один момент времени задачу выполняет только один процесс.
Выводы
Главный вывод для нас оказался простым: правильная модель важнее конкретного инструмента.
В какой‑то момент мы перестали искать «правильный лок» и начали разбираться, что именно хотим контролировать. Оказалось, что нам нужно было синхронизировать не процессы, а окна расписания. После этого найти решение стало намного проще.Теги:• python
• распределенные системы
• блокировки
• distributed systems
• locksХабы:• Блог компании Точка Банк
• Python
• Веб-разработка
• Качество кода
• Программирование
Получайте больше инсайтов о систематизации бизнеса
Подписывайтесь на Telegram-канал Business Operations — ежедневные материалы о бизнес-процессах, операционном управлении и повышении эффективности
💬 Подписаться на канал→ Оригинальная статья