redb.Route.Crone
redb.Route.Crone


redb ecosystem

Есть таблица GPS-точек. Транспорт шлёт координаты каждые несколько секунд, за сутки набегают миллионы строк, и таблица секционирована по времени. Значит, кто-то должен заранее создавать партицию на завтра и отцеплять партиции старше девяноста дней. Иначе в одну прекрасную ночь вставка упадёт с no partition of relation "gps_points" found for row, и это будет ровно в тот момент, когда никто не смотрит.

redb.Tsak scheduler
redb.Tsak scheduler

Задача старая как мир, и обычно решается одним из двух способов: pg_cron внутри базы или IHostedService с PeriodicTimer в приложении. У первого проблема с наблюдаемостью — джоб живёт в базе, а логи приложения о нём ничего не знают. У второго проблема с тем, что вокруг PeriodicTimer очень быстро нарастает своя маленькая инфраструктура: retry, логирование, «а что если предыдущий запуск ещё идёт», «а что если нод три».

Я покажу третий способ — маршрут. Дальше будет разбор коннектора redb.Route.Quartz, вызов PostgreSQL-функции через SQL-коннектор, сплиттер с изоляцией ошибок и честный разговор про кластер, потому что именно там всё интересное и начинается.

Весь код в статье — на строковых URI. Fluent-билдеры в redb.Route есть, но URI читается без знания API, и его можно скопировать в конфиг.

Серия про экосистему redb и redb.Route. Это продолжение цикла, свежие статьи — сверху:

Полный список — в профиле. Исходники: 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
  &param.keep_days=90

procedureParams — это объявление параметров в формате направление:имя:тип, порядок объявления и есть порядок аргументов в вызове. Направления три: INOUTINOUT; значения OUT после выполнения возвращаются обратно в заголовки сообщения под своими именами.

Значение параметра ищется по цепочке: сначала явный param.имя из URI, потом заголовок сообщения с таким же именем, потом тело, если оно словарь. Здесь keep_days задан константой прямо в URI, а tbl придёт из заголовка — его выставит сплиттер.

Для простых случаев есть более короткий путь — mode=Execute (он же дефолтный) с плейсхолдерами @имя прямо в тексте запроса:

sql:SELECT maintain_partitions(@tbl, @keep_days)
  ?dataSource=#pg
  &outputType=Scalar
  &param.tbl=${header.tbl}
  &param.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"
                + "&param.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.cs
if (!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, а вся информация о срабатывании лежит в свойствах:

Свойство

Что внутри

CamelQuartzFireTime

когда джоб фактически сработал

CamelQuartzScheduledFireTime

когда должен был сработать по расписанию

CamelQuartzNextFireTime

когда сработает в следующий раз

CamelQuartzPreviousFireTime

когда срабатывал в прошлый раз

CamelCronSchedule / CamelCronName / CamelCronGroup

само выражение, имя и группа джоба

Разница между 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-имя), threadsmisfireInstructionstatefulrecoverableJobdurableJobdeleteJobpauseJobstartAtendAtcustomCalendartriggerStartDelayprefixJobNameWithEndpointId.

У qtimer: вместо schedule — perioddelayfixedRaterepeatCount, остальное то же самое, кроме таймзоны (простому триггеру она не нужна).

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.