Наш пакетный планировщик на самом деле load balancer
Каждую ночь Norman балансирует на одних и тех же воркерах двух очень разных «жильцов»: терпеливые, сетевые банковские синхронизации и голодные, вычислительно тяжёлые AI-джобы, которые читают чеки и категоризируют постоянный поток транзакций. Вот как мы относимся к планировщику как к load balancer, измеряем давление очередей по каждой нагрузке и куда ведём это дальше: раскладываем тяжёлую AI-работу по живой нагрузке, а не по вручную выбранной минуте крона.
- Категория
- Общее
- Обновлено
- Автор
- Stan Kharlap
Большую часть того, что делает Norman, ты видишь. Открываешь приложение: транзакция категоризирована, чек сопоставлен, декларация по НДС готова. Чего ты не видишь, так это ночной смены. Пока Германия спит, флот фоновых воркеров балансирует на одном и том же железе двух очень разных «жильцов»: терпеливые, сетевые банковские синхронизации по тысячам подключений и голодные, вычислительно тяжёлые AI-джобы, которые читают чеки, категоризируют постоянный поток новых транзакций и переигрывают вчерашние эвалы агентов. Ничего гламурного в этом нет, и ни в чём из этого не оказалось тех сложных задач, которых я ожидал.
Сложная задача в том, что эти двое жильцов хотят одну и ту же машину в одно и то же время и нагружают её совершенно по-разному. Фоновый планировщик выглядит как список записей «запусти это ночью». На практике это load balancer в шляпе крона: его настоящая работа не дать одной нагрузке, и почти всегда это AI-нагрузка, уморить остальные голодом. Скучные решения о том, что работает рядом с чем, и составляют бо́льшую часть того, что держит утро спокойным.
Окно это общий ресурс
Вот в какой отказ мы вошли. Наши самые тяжёлые повторяющиеся джобы съехались все в одно узкое окно в ранние часы, потому что это очевидное место для «ночной» работы. Внутри этого одного окна расписание выглядело так:
update_all_transactions, который разворачивает open-banking-синхронизацию по каждому подключённому счёту, тысячам счетов, и каждый это медленный вызов к внешнему агрегатору.update_accounts, обновление сальдо, ещё внешние вызовы.categorize_uncategorized_transactions, тяжёлый прогон AI и OCR постобработки, прямо поверх остального.sync_all_active_integrations, закрывающий окно.
Эти двое пиковали друг против друга, и они противоположны. Open-banking-развёртка терпеливая и медленная, она проводит время в ожидании чужого API. AI-прогон нетерпеливый и голодный: ему нужны воркеры и мощность модели сейчас. В одном окне они не складываются, а мешают друг другу. Синхро-джобы держат воркеров заложниками в ожидании ввода-вывода, а AI-джобы стоят в очереди позади них.
Решение было почти неловким в своей простоте, и я вернусь к тому, почему. Но ты не сделаешь его безопасно, пока работа не разделена по типу, и вот тут-то и лежит настоящий дизайн.
Сначала раздели работу по тому, чего она ждёт
Самое полезное, что мы сделали, это перестали относиться к «фоновому джобу» как к одной категории. Джоб, который ждёт банковский API, и джоб, который ждёт языковую модель, падают по-разному, ретраятся по-разному и морят друг друга голодом, если делят один пул воркеров. Поэтому они его не делят.
Работа маршрутизируется по отдельным очередям в зависимости от того, чего она ждёт, у каждой свой пул воркеров:
# Маршрутизируем каждую задачу в очередь по тому, чего она ждёт (иллюстративно):
app.conf.task_routes = {
# пользовательская AI / OCR / работа с документами
"categorize_transaction": {"queue": "transaction_postprocess"},
"extract_data_from_attachment": {"queue": "transaction_postprocess"},
# вызовы open-banking API: сетевые и медленные
"update_transactions_by_bank_account": {"queue": "transaction_sync"},
# фоновая подготовка агентов vs. видимые пользователю действия агентов
"run_evals": {"queue": "agent_batch"},
"submit_approved_run": {"queue": "agent_realtime"},
# ...
}
Различия, которые важны, это профили нагрузки, а не имена. transaction_sync полон джобов, которые для нас дёшевы, а медленны из-за кого-то другого, так что его пул держит много джобов, в основном простаивающих на сокете. transaction_postprocess наоборот: дорогой в пересчёте на джоб, и ты не хочешь, чтобы наводнение приземлилось разом. Две очереди агентов кодируют скорее разделение по приоритету, чем по ресурсу: agent_batch это подготовка, которую никто не ждёт, а agent_realtime это момент, когда человек нажал «Отправить» и смотрит на спиннер. Развести их по разным очередям означает, что ночной прогон эвалов никогда не встанет перед декларацией, которую человек пытается отправить.
Как только полосы существуют, планирование превращается в вопрос о том, какие полосы ты зажигаешь вместе.
Потом не планируй две тяжёлые вещи одновременно
С разделённой работой само изменение планировщика было в пару строк крона. Мы увели трио AI и агентов из банковского окна целиком:
"categorize_uncategorized_transactions": {
# Уведён из open-banking-окна, чтобы тяжёлый прогон AI/OCR постобработки
# не пиковал одновременно с развёрткой синхронизации; теперь он идёт
# в более тихом слоте позже ночью.
"schedule": crontab(hour=..., minute=...),
"options": {"expires": ...}, # щедрое окно, час-два
},
"promote_findings": {
"schedule": crontab(hour=..., minute=...), # прямо перед run_evals
"options": {"expires": ...},
},
"run_evals": {
"schedule": crontab(hour=..., minute=...), # ночной повтор эвалов
"options": {"expires": ...},
},
Прогон AI теперь идёт в более тихом слоте позже ночью, после ночного служебного аудита и с запасом до утренней отправки счетов. Банковское окно оставляет свой ранний слот себе. В самих джобах ничего не поменялось, они просто перестали сталкиваться.
Скажу честно: это не хитрое изменение. Это правка расписания. Она доступна нам как правка расписания только потому, что менее заметная работа ниже, разделение очередей и семантика ретраев, сделала каждый из этих джобов безопасным для переноса на пару часов. Когда джобы сплетены, перенести один это оценка рисков; когда они изолированы и идемпотентны, это однострочный диф. Бо́льшая часть ценности в том, чтобы заработать право делать скучные изменения.
expires: опоздавший джоб тебе не нужен
Присмотрись к этим записям ещё раз, и увидишь, что каждая несёт окно expires, измеряемое в часах. Это не украшение. Это ответ на конкретный вопрос: что должно случиться, если запланированный джоб отправлен, а воркеры слишком заняты, чтобы подхватить его вовремя?
Неправильный ответ это «запусти его, когда освободится воркер». Ночной прогон категоризации полезен, когда он срабатывает, в ранние часы. Если пул был перегружен и тот же прогон в итоге подхватывается в середине дня, он теперь конкурирует с живым пользовательским трафиком за работу, которую более поздний прогон всё равно переделает. Устаревший пакетный джоб, повторённый посреди дня, хуже, чем пропущенный пакетный джоб.
Поэтому beat-записи истекают. Если запланированный джоб не стартовал внутри своего окна, брокер его отбрасывает, а следующий запланированный прогон разбирается с накопившимся. expires превращает «когда-нибудь» в «сейчас или никогда», а это ровно тот контракт, который нужен для периодической работы. Окна нарочно щедрые, чтобы нормальная смена спокойно помещалась, и только по-настоящему застрявший пул вообще запускает отбрасывание.
Меряем нагрузку непрерывно
Load balancer, который не видит нагрузку, это просто статический роутер. Поэтому в планировщике на коротком интервале работает зонд, который читает реальную глубину каждой очереди по нагрузке, а не одно число «занята ли коробка»:
QUEUE_SIZE_THRESHOLD = ... # пара сотен ожидающих сообщений
MONITORED_QUEUES = [
"transaction_sync", # развёртка банковской синхронизации
"transaction_postprocess", # AI / OCR
"agent_batch", "agent_realtime",
]
@app.task(name="monitoring.check_queue_health")
def check_queue_health() -> None:
"""Алертит, когда какая-то очередь по нагрузке распухает."""
for name in MONITORED_QUEUES:
if queue_depth(name) > QUEUE_SIZE_THRESHOLD:
logger.critical("Очередь %s распухает", name)
# ...и отдельный присмотр за dead-letter-очередью
Тут важны две вещи. Во-первых, он читает каждую нагрузку отдельно, так что мы никогда не видим просто «система занята», мы видим, какой жилец под давлением: полоса банковской синхронизации, залипшая на медленном агрегаторе, или полоса AI и OCR, перемалывающая накопившееся. По одному графику CPU их не различить, а требуют они противоположных реакций. Во-вторых, он следит за dead-letter-очередью, куда задачи попадают, провалившись на раз больше положенного, потому что растущий счётчик dead-letter значит, что что-то не просто медленно, а сломано.
Сегодня работа зонда это пейджить человека. Но он же и сенсор, который нужен более умному планировщику. Нельзя балансировать нагрузку, которую не можешь измерить, а это та часть, что её меряет, по каждой нагрузке, непрерывно.
Как минимум один раз значит идемпотентно или никак
Другая вещь, которая делает джобы безопасными для перепланирования, это что происходит, когда один умирает на полпути. Мы гоняем воркеров с двумя настройками, которые меняют одну проблему на проблему получше:
task_acks_late=True,
task_reject_on_worker_lost=True,
acks_late означает, что джоб подтверждается только после завершения, а не когда его подхватили, так что воркер, упавший на полпути, не теряет джоб молча. reject_on_worker_lost кладёт этот джоб обратно в очередь. Вместе они покупают тебе доставку как минимум один раз: при падении ничего не теряется. Цена в том, что джоб иногда может выполниться дважды, один раз на воркере, который умер, и один раз на его замене.
А значит, каждый джоб должен быть безопасен для двойного запуска, и этот инвариант тихо формирует всю систему. Слой AI-агентов штампует каждый прогон ключом idempotency_key, так что повторный вызов модели разрешается в уже существующий прогон, а не в свежий дубликат, на десятках тысяч прогонов агентов, что мы записали. Путь банковской синхронизации дедуплицирует на уровне счёта, блокировкой на счёт при каждой отправке, так что агрессивное планирование никогда не превращается в агрессивную двойную работу. Смысл доставки как минимум один раз в том, что ты перестаёшь предотвращать дубликаты на очереди и вместо этого делаешь их безвредными на джобе.
Прими, что он застрянет, и построй метлу
Последняя категория запланированных джобов существует чисто потому, что остальные будут вести себя плохо. Распределённая пакетная работа не падает чисто, она застревает. Воркер умирает между двумя коммитами, внешний вызов зависает, блокировку держит процесс, которого больше нет. Так что добрая часть расписания это уборка:
recover_stuck_syncsидёт регулярным проходом и расклинивает банковские счета, застрявшие посреди синхронизации.requeue_stuck_runsвосстанавливает прогоны агентов, чья воркер-задача была потеряна.update_all_transactionsидёт несколько раз в день, не потому что банки меняются так часто, а как страховочная сетка, переподталкивающая счета, что предыдущий проход оставил позади.
Ни один из них ничего не делает в хорошую ночь. Это метла в шкафу. Но допущение за ними, что любой отдельный прогон может упасть и упадёт, ровно и держит счастливый путь простым: ни одному джобу не нужно быть пуленепробиваемым, потому что придёт уборщик и переподтолкнёт всё, что тот оставил в плохом состоянии.
К планировщику, который балансирует сам себя
Всё до этого нарочно статично. Очереди фиксированы, минуты крона выбраны вручную, и когда AI-прогон столкнулся с окном банковской синхронизации, человек это заметил и перенёс. Это работает, и скучно-и-верно всякий раз бьёт хитро-и-хрупко. Но отсюда уже видна форма следующего шага.
У нас уже есть обе половины, которые нужны адаптивному балансировщику. Мы умеем мерять нагрузку: зонд глубины очередей читает живое давление по каждой нагрузке и следит за dead-letter-очередью на предмет поломок. И мы умеем раскладывать нагрузку: таблица маршрутизации решает, где идёт джоб, расписание решает, когда. Не хватает контроллера посередине, той части, что читает сигнал и двигает работу, не дожидаясь, пока человек заметит.
Направление, которое мы исследуем, это планировщик, который относится к AI- и OCR-джобам как к эластичным жильцам. Вместо фиксированного «запусти прогон категоризации в эту минуту» он спрашивал бы в момент диспатча, где есть запас: придержать тяжёлую вычислительную работу, пока пикует развёртка банковской синхронизации, и отпустить её в следующий тихий слот. Это то самое решение, что мы приняли вручную в этом месяце, принятое непрерывно из живой нагрузки. Сложное в адаптивном load balancer это никогда не политика раскладки; это возможность двигать работу, ничего не ломая, и это основание уже заложено. Каждый джоб идемпотентен, так что контроллер волен перекладывать свободно; у каждого джоба есть окно expires, так что неудачно положенный будет отброшен, а не выполнен устаревшим; каждая нагрузка изолирована, так что перенос одной никогда не морит другую голодом. Умная часть маленькая. Скучная часть под ней и есть весь продукт.
Скучные части и есть продукт
Весь планировщик это пара сотен строк Python: словарь крон-записей, таблица маршрутизации и зонд, который читает глубину очередей. Изменение, которое породило этот пост, перенесло горстку этих записей в более тихое окно. Если ты искал AI в AI-продукте для бухгалтерии: модель, читающая твои чеки, здесь всего лишь самый громкий жилец, а не сложная часть.
Сложная часть это слой вокруг неё: успеет ли эта модель отработать до твоего утра, встанет ли декларация, которую ты пытаешься отправить, позади ночного эвала, и не переотправит ли тихо падение воркера в неудачный момент то, что не должно. Интересная инженерия в системе вроде нашей это редко вызов модели. Это негламурный балансировщик, который решает, когда модель идёт, рядом с чем она идёт и что происходит, когда не идёт. Мы тратим много времени на этот балансировщик именно затем, чтобы тебе никогда не приходилось о нём думать.
Norman берет операционную финансовую работу на себя
От invoicing до bookkeeping: Norman организует повторяющиеся финансовые процессы так, чтобы вы успевали к дедлайнам с меньшим объемом ручной работы.