redb ecosystem
Есть таблица GPS-точек. Транспорт шлёт координаты каждые несколько секунд, за сутки набегают миллионы строк, и таблица секционирована по времени. Значит, кто-то должен заранее создавать партицию на завтра и отцеплять партиции старше девяноста дней. Иначе в одну прекрасную ночь вставка упадёт с no partition of relation "gps_points" found for row, и это будет ровно в тот момент, когда никто не смотрит.
Задача старая как мир, и обычно решается одним из двух способов: pg_cron внутри базы или IHostedService с PeriodicTimer в приложении. У первого проблема с наблюдаемостью — джоб живёт в базе, а логи приложения о нём ничего не знают. У второго проблема с тем, что вокруг PeriodicTimer очень быстро нарастает своя маленькая инфраструктура: retry, логирование, «а что если предыдущий запуск ещё идёт», «а что если нод три».
Я покажу третий способ — маршрут. Дальше будет разбор коннектора redb.Route.Quartz, вызов PostgreSQL-функции через SQL-коннектор, сплиттер с изоляцией ошибок и честный разговор про кластер, потому что именно там всё интересное и начинается.
Весь код в статье — на строковых URI. Fluent-билдеры в redb.Route есть, но URI читается без знания API, и его можно скопировать в конфиг.
Серия про экосистему redb и redb.Route. Это продолжение цикла, свежие статьи — сверху:
redb.Route — коннектор RabbitMQ: RPC, конкурирующие консьюмеры и dead-letter. Уходим от MassTransit
redb.Route: два маршрута за вечер — от отладочного воркера до энтерпрайза на Tsak
redb.Route — уходим от MassTransit, идём к Apache Camel: Kafka, Scatter‑Gather и транзакции
Полный список — в профиле. Исходники: github.com/redbase-app/redb-route. Про саму БД: redb.ru.
Сначала SQL
Никакой магии на этом уровне не будет, поэтому начнём с самого честного места — с функции в базе. Она принимает имя таблицы и срок хранения, создаёт партицию на завтра, если её ещё нет, и отцепляет всё старше срока:
CREATE OR REPLACE FUNCTION maintain_partitions(tbl text, keep_days int)RETURNS text AS $$DECLARE next_day date := (now() + interval '1 day')::date; part_name text := format('%s_%s', tbl, to_char(next_day, 'YYYYMMDD')); cutoff date := (now() - make_interval(days => keep_days))::date; dropped int := 0; old_part text;BEGIN IF to_regclass(part_name) IS NULL THEN EXECUTE format( 'CREATE TABLE %I PARTITION OF %I FOR VALUES FROM (%L) TO (%L)', part_name, tbl, next_day, next_day + 1); END IF; FOR old_part IN SELECT c.relname FROM pg_class c JOIN pg_inherits i ON i.inhrelid = c.oid JOIN pg_class p ON p.oid = i.inhparent WHERE p.relname = tbl AND c.relname < format('%s_%s', tbl, to_char(cutoff, 'YYYYMMDD')) LOOP EXECUTE format('DROP TABLE %I', old_part); dropped := dropped + 1; END LOOP; RETURN format('%s: +1 partition, -%s dropped', tbl, dropped);END;$$ LANGUAGE plpgsql;
Функция возвращает строку, чтобы её было видно в логе. Вот это «чтобы было видно в логе» — единственная уступка удобству, всё остальное здесь обычный plpgsql, который вы бы написали в любом случае.
Тик: cron://
Коннектор redb.Route.Quartz даёт две схемы. Первая — cron:, обычный quartz-овский триггер с cron-выражением:
cron://[группа/]имяДжоба?schedule=<cron-выражение>&<опции>
Вторая — qtimer:, простой периодический триггер, если cron не нужен:
qtimer://[группа/]имяДжоба?period=5000&delay=1000&fixedRate=true
Схемы quartz: намеренно нет. Это не забывчивость: настройка планировщика (где хранятся джобы, кластер это или одна нода, какой пул потоков) — свойство хоста, а не маршрута. В URI маршрута попадает только расписание. Если бы существовала схема quartz:, в неё немедленно начали бы прорастать настройки job store, и маршрут перестал бы быть переносимым.
Cron-выражение — квартцевское, шестипольное, с секундами. Обслуживание партиций поставим на 02:30:
From("cron://maintenance/gps-partitions?schedule=0 30 2 * * ?") .RouteId("gps-partitions")
Выражение валидируется сразу при создании эндпоинта, а не в момент первого срабатывания. Опечатка в расписании — это ArgumentException на старте приложения, а не тишина до трёх часов ночи.
Регистрация компонента — одна строка в точке входа модуля:
context.AddComponent(new CronComponent());
Вызов функции: sql:
SQL-коннектор — одна схема sql:, режим выбирается параметром mode. Для вызова PostgreSQL-функции есть mode=Procedure с флагом asFunction=true — тогда коннектор соберёт SELECT maintain_partitions(@p1, @p2) и выполнит скаляром, положив результат в тело сообщения:
sql:maintain_partitions ?mode=Procedure &dataSource=#pg &procedureName=maintain_partitions &asFunction=true &procedureParams=IN:tbl:String,IN:keep_days:Int32 ¶m.keep_days=90
procedureParams — это объявление параметров в формате направление:имя:тип, порядок объявления и есть порядок аргументов в вызове. Направления три: IN, OUT, INOUT; значения OUT после выполнения возвращаются обратно в заголовки сообщения под своими именами.
Значение параметра ищется по цепочке: сначала явный param.имя из URI, потом заголовок сообщения с таким же именем, потом тело, если оно словарь. Здесь keep_days задан константой прямо в URI, а tbl придёт из заголовка — его выставит сплиттер.
Для простых случаев есть более короткий путь — mode=Execute (он же дефолтный) с плейсхолдерами @имя прямо в тексте запроса:
sql:SELECT maintain_partitions(@tbl, @keep_days) ?dataSource=#pg &outputType=Scalar ¶m.tbl=${header.tbl} ¶m.keep_days=90
Обратите внимание на ${header.tbl} — это выражение, оно резолвится в рантайме из заголовка сообщения. Плейсхолдеры в SQL — только @имя, двоеточие не поддерживается, и неподставленный параметр молча становится NULL, так что имена лучше не путать.
Оба варианта рабочие. Дальше в статье я использую mode=Procedure, потому что он показательнее.
Сплиттер: три таблицы, одна упавшая не роняет остальные
Таблиц с временными партициями обычно не одна. У нас их три: точки, треки и события. Наивно было бы написать цикл внутри процессора, но тогда придётся руками решать, что делать, если вторая таблица упала: прервать всё или продолжить, и как потом понять, что именно не отработало.
Это ровно задача паттерна Splitter из EIP. Сообщение со списком разбивается на сообщения по элементу, каждое идёт по своей ветке, и ветки можно обрабатывать параллельно:
From("cron://maintenance/gps-partitions?schedule=0 30 2 * * ?") .RouteId("gps-partitions") .Process(e => e.In.Body = new[] { "gps_points", "gps_tracks", "gps_events" }) .Split(Body()) .ParallelProcessing() .MaxParallelism(2) .SetHeader("tbl", Body()) .DoTry() .To("sql:maintain_partitions" + "?mode=Procedure" + "&dataSource=#pg" + "&procedureName=maintain_partitions" + "&asFunction=true" + "&procedureParams=IN:tbl:String,IN:keep_days:Int32" + "¶m.keep_days=90") .Log("[PART] ${body}") .DoCatch<Exception>() .Log("[PART] ${header.tbl}: ${exception.Message}", LogLevel.Error) .End() .EndSplit() .Process(Summary);
Двенадцать строк, и в них уже есть всё, что обычно дописывают руками неделю спустя. MaxParallelism(2) — две таблицы обслуживаются одновременно, третья ждёт свободного слота; DROP TABLE берёт ACCESS EXCLUSIVE, и заваливать базу параллельными блокировками смысла нет. DoTry/DoCatch стоят внутри сплита, поэтому исключение изолировано в своей ветке: упавшая gps_tracks не отменит уже отработавшую gps_points и не помешает gps_events. После EndSplit управление приходит в Summary, где можно посчитать, сколько веток отработало, и решить, звать ли дежурного.
Это не выдуманный ради статьи приём. Ровно такая конструкция крутится в проде — джоб синхронизации точек отгрузки из SAP разбивает список точек и обрабатывает по три параллельно, потому что одна недоступная точка не должна ронять весь тик:
From("timer://tsum-points?period=180000&delay=60000") .RouteId("tsum-points-timer") .ProcessWithRedb(PreloadContextAsync) .Split(Body()) .ParallelProcessing() .MaxParallelism(3) .SetHeader("ShippingPoint", Body()) .DoTry() .To(sqlTo) .Process(DeserializeXml) .ProcessWithRedb(ProcessPointsAsync) .DoCatch<Exception>() .Process(AddPointSyncError) .Log(LogLevel.Error) .Message("[TSUM-PT] SP=${header.ShippingPoint} failed, skipping: ${exception.Message}") .EndLog() .End() .EndSplit() .Process(BuildPointSyncSummary);
Сплиттер — не единственный EIP, который стыкуется с планировщиком. По расписанию естественно ложатся Content-Based Router (в будни одно, в выходные другое), Throttler (не долбить внешний API чаще N раз в секунду), Aggregator (собрать результаты веток в один отчёт), Idempotent Consumer (о нём ниже) и Dead Letter Channel. В redb.Route реализованы почти все паттерны каталога Хоупа и Вульфа — Splitter, Aggregator, Resequencer, Multicast, Recipient List, Dynamic Router, Wire Tap, Content Enricher, Claim Check, Saga, Scatter-Gather, Circuit Breaker, Load Balancer, Transactional Client и остальные. Планировщик здесь просто источник, а не отдельный мир со своими правилами.
Что происходит, когда сервер лежал в 02:30
Вот тут начинается то, ради чего вообще нужен Quartz, а не PeriodicTimer.
Джоб не выстрелил, потому что нода была в дауне или деплой затянулся. Что делать, когда планировщик поднялся в 02:47? Ответ зависит от джоба, и это не философский вопрос: для обслуживания партиций пропуск — катастрофа, партицию на завтра надо создать хоть в 02:47, хоть в 06:00. А для джоба «разослать утренний отчёт» запуск в полдень — хуже, чем ничего.
Quartz называет это misfire, и коннектор пробрасывает политику прямо в URI:
cron://maintenance/gps-partitions?schedule=0 30 2 * * ?&misfireInstruction=CronFireOnceNow
CronFireOnceNow — догнать, выполнить один раз и вернуться в расписание. CronDoNothing — пропустить, ждать следующего по расписанию. Для simple-триггеров (qtimer:) политик больше — пять штук, они отличаются тем, что делать с накопившимся счётчиком повторов. Но выбор всегда сводится к одному вопросу: пропущенный запуск нужно догнать или он уже протух?
Дефолт коннектора для qtimer: выбран по здравому смыслу: при fixedRate=true — догнать (вы просили фиксированную частоту), иначе — перепланировать со следующего тика.
Что происходит, когда предыдущий запуск ещё идёт
Партиций накопилось много, DROP TABLE ждёт блокировку, джоб висит. Наступает следующее срабатывание. Что делать?
Стандартный ответ Quartz — атрибут [DisallowConcurrentExecution] на классе джоба. Коннектор пошёл другим путём: конкурентность контролирует семафор самого консьюмера, и если все потоки заняты, срабатывание молча пропускается:
// QuartzConsumerBase.csif (!await _semaphore.WaitAsync(0).ConfigureAwait(false)) return; // все потоки заняты — пропускаем это срабатывание
Размер семафора задаётся в URI параметром threads (по умолчанию 1). Почему не атрибутом: [DisallowConcurrentExecution] — это либо один запуск, либо никакого контроля, промежуточных значений нет. А семафор позволяет сказать «до трёх параллельных запусков этого джоба» и при этом остаётся дружелюбным к кластеру, где ограничение живёт на уровне job store, а не атрибута класса.
Если всё-таки нужна квартцевская семантика — есть флаг stateful=true, он переключает джоб на класс с [DisallowConcurrentExecution] и [PersistJobDataAfterExecution].
Ещё одна деталь про остановку: когда маршрут гасится, коннектор снимает триггер и ждёт завершения уже выполняющихся запусков — до тридцати секунд. Джоб, который в этот момент дропает партицию, не будет прерван на полпути. А если джоб всё-таки сработал, а маршрута уже нет (например, модуль выгрузили) — джоб при запуске обнаруживает, что его консьюмер мёртв, и удаляет себя из планировщика сам, не оставляя мусора.
Три ноды
Самый частый вопрос про cron в распределённом приложении: если нод три, джоб выстрелит трижды?
Выстрелит — если каждая нода держит свой планировщик в памяти. Именно так работает fallback коннектора: не нашёл IScheduler в контексте — создал свой, с RAM-хранилищем, уникальный для этого контекста. Для локальной разработки этого достаточно, для трёх нод — нет.
Правильный ответ — один планировщик на кластер, а точнее одно общее хранилище джобов. Quartz умеет это из коробки через AdoJobStore, и коннектор специально не мешает: он не создаёт свой планировщик, если готовый уже лежит в контексте маршрута. Хост подкладывает туда кластерный, и никаких изменений в маршруте не требуется — URI остаётся тем же.
Вот тут и вылезает наружу вся кухня, без которой обычно обещают обойтись. AdoJobStore — это набор таблиц QRTZ_* в вашей базе. QRTZ_TRIGGERS хранит следующее время срабатывания, QRTZ_FIRED_TRIGGERS — кто что сейчас выполняет, QRTZ_LOCKS — строки-мьютексы. Механизм «только одна нода выполнит джоб» — это не хитрый консенсус, а SELECT ... FOR UPDATE по строке в QRTZ_LOCKS: кто первым взял блокировку, тот и забрал триггер. Каждая нода периодически отмечается в QRTZ_SCHEDULER_STATE, и если нода перестала отмечаться, её незавершённые джобы подхватывает другая — но только те, у которых стоит recoverableJob=true. По умолчанию флаг выключен: перезапускать джоб, о котором вы не знаете, идемпотентен он или нет, — плохая идея.
Схема таблиц создаётся хостом на старте, строка подключения и диалект берутся из конфигурации базы приложения. То есть DDL вы не пишете, но таблицы — самые обычные, лежат рядом с вашими, видны в любом клиенте, и когда что-то пойдёт не так, вы залезете туда обычным SELECT и увидите, какой триггер завис и на какой ноде.
И раз уж мы включили recoverableJob: перезапуск после падения ноды означает, что джоб может выполниться дважды. Для нашей функции обслуживания партиций это безопасно — она написана идемпотентно (IF to_regclass(...) IS NULL), и это не случайность, а требование. Если бы джоб был неидемпотентным — скажем, начислял бонусы, — перед ним нужно ставить Idempotent Consumer, и в redb.Route для него есть репозиторий поверх SQL с уникальным индексом. Уникальный индекс, а не «умный кеш»: в кластере от двойного выполнения спасает только он.
Что кладётся в сообщение
Планировщик — источник без тела сообщения. Тело null, паттерн InOnly, а вся информация о срабатывании лежит в свойствах:
|
Свойство |
Что внутри |
|---|---|
|
|
когда джоб фактически сработал |
|
|
когда должен был сработать по расписанию |
|
|
когда сработает в следующий раз |
|
|
когда срабатывал в прошлый раз |
|
|
само выражение, имя и группа джоба |
Разница между FireTime и ScheduledFireTime — это ровно тот самый misfire в цифрах. Если они разошлись на семнадцать минут, значит, джоб догоняли.
Имена в стиле Camel — не ностальгия. redb.Route сознательно сохраняет номенклатуру Apache Camel там, где семантика совпадает, чтобы человек, приходящий из Java-интеграций, читал заголовки без словаря.
Два джоба из прода
Чтобы не выглядело как статья про сферический cron в вакууме — вот два маршрута, которые каждую ночь работают в системе управления транспортом.
Бэкап базы в три часа:
From("cron://tsum-backup?schedule=0 0 3 * * ?") .RouteId("tsum-backup-cron") .ProcessWithRedb(RunBackupAsync);
Чистка мёртвых маршрутов в четыре:
From("cron://tsum-cleanup?schedule=0 0 4 * * ?") .RouteId("tsum-cleanup-cron") .ProcessWithRedb(CleanupDeadRoutesAsync);
Ничего эффектного, и это хорошо: расписание в URI, логика в процессоре, ретраи и логирование — от фреймворка. Обратите внимание, что расписание захардкожено в URI, а не вынесено в конфиг — так тоже можно, и первое время так и живут. Когда понадобится менять расписание без пересборки, URI собирается из конфига обычной конкатенацией, потому что это просто строка.
Полный список опций
Чтобы не разворачивать справочник на пол-статьи — всё, что понимает cron::
schedule (обязательный), timeZone (IANA-имя), threads, misfireInstruction, stateful, recoverableJob, durableJob, deleteJob, pauseJob, startAt, endAt, customCalendar, triggerStartDelay, prefixJobNameWithEndpointId.
У qtimer: вместо schedule — period, delay, fixedRate, repeatCount, остальное то же самое, кроме таймзоны (простому триггеру она не нужна).
customCalendar стоит отдельного упоминания: это квартцевский календарь исключений, зарегистрированный в контексте по имени. Через него делаются «кроме праздников» и «только в рабочие дни» — то, что в cron-выражении не выражается никак.
Итог
Планировщик в redb.Route — это источник сообщений, а не отдельная подсистема со своими правилами. Джоб — это маршрут, а значит, ему доступно всё, что доступно любому маршруту: сплиттер, обработка исключений, транзакции, ретраи, метрики. Триггер описывается одной строкой URI, и в этой строке нет ничего про инфраструктуру — только расписание.
Инфраструктура при этом никуда не спрятана. Кластер работает на таблицах QRTZ_* и блокировке строки в базе, идемпотентность в кластере обеспечивается уникальным индексом, а обслуживание партиций — обычной plpgsql-функцией, которую вы написали и можете прочитать. Фреймворк здесь избавляет от связующего кода, а не от понимания того, что происходит в базе.
Код: github.com/redbase-app/redb-route · сайт: redb.ru
If this was useful — a ⭐ on GitHub helps others find it.
ссылка на оригинал статьи https://habr.com/ru/articles/1061558/