Когда‑то наш почтовый загрузчик был обычным методом с @Scheduled. Раз в 30 минут он подключался к ящикам, искал новые вложения и запускал их обработку. Решение было простым, понятным и долгое время вполне рабочим.
Через несколько итераций вокруг того же загрузчика уже существовали IMAP IDLE listener'ы, PostgreSQL leases, heartbeat экземпляров приложения, перебалансировка mailbox'ов между pod'ами, UID checkpoints, постоянная идемпотентность и отдельная обработка FolderClosedException. В какой‑то момент мы даже попробовали держать несколько IMAP‑соединений к одному ящику — а затем сознательно удалили эту часть архитектуры.
Это история не о том, как мы «заменили polling на push» одной настройкой. Она о том, как безобидный интеграционный адаптер постепенно превратился в маленькую распределённую систему. И о границе параллелизма, которую мы сначала провели не там.

Что вообще делал этот сервис
Наш продукт обрабатывает прайс‑листы поставщиков. В них приходят товарные позиции, цены и сведения об остатках. Для значительной части поставщиков почта была одним из основных каналов обновления: поставщик отправляет письмо с файлом, а дальше система должна сама понять, к какой конфигурации относится вложение, и запустить обновление.
При этом mailbox и прайс‑лист — не одно и то же. Один физический почтовый ящик мог обслуживать несколько логических конфигураций. Они различались отправителем, темой письма, ожидаемым именем вложения и правилами дальнейшей обработки.
Поэтому получение письма было только началом. Подходящее вложение нужно было сохранить, при необходимости извлечь из архива, определить формат, распарсить и классифицировать данные, записать результат, а затем сформировать события и уведомления об исходе обработки.
В упрощённом виде путь выглядел так:

Свежесть здесь имела два определения.
Первое очевидно: актуальные цены и остатки хотелось получать сразу, а не через случайный интервал от нуля до 30 минут. Второе проявилось не сразу. За одно такое окно один поставщик мог несколько раз прислать прайс на актуализацию. Для нас была важна не только последняя версия файла, но наблюдаемая история обработки каждого входящего обновления.
Polling: сначала это было хорошее решение
Историческая реализация, если убрать бизнес‑детали, выглядела примерно так:
@Scheduled(fixedRate = 30 * 60 * 1000) public void scheduleMailboxProcessing() { loadEnabledPriceLists() .forEach(this::checkMailbox); }
У такого подхода было много достоинств.
Во‑первых, он был предельно прозрачен. Есть расписание, есть один запуск, есть понятный набор логов. Если почтовый сервер временно недоступен, следующий цикл естественным образом превращается в повторную попытку.
Во‑вторых, не нужно управлять долгоживущими соединениями. Нет отдельного lifecycle для listener'ов, не нужно восстанавливать IDLE после сетевого разрыва и синхронизировать состояние при изменении конфигурации.
В‑третьих, polling хорошо ограничивает точки отказа. Сломался один проход — через 30 минут начнётся следующий. Для первой версии интеграции это вполне оптимальный обмен сложности на задержку.
Проблема была не в том, что polling «плохой». Мы просто эволюционно переросли гарантии, которые давала конкретная реализация.
Как три письма превращались в одно событие
Каждый проход загружал ограниченное окно последних сообщений и сортировал их по времени отправки — от новых к старым. Затем алгоритм последовательно проверял тему, отправителя и вложение. На первом подходящем письме поиск завершался.
После успешной обработки timestamp выбранного сообщения сохранялся как новый checkpoint конфигурации. На следующем проходе более старые письма уже считались обработанным прошлым.
Рассмотрим одно 30-минутное окно:
10:05 письмо A с подходящим прайсом 10:12 письмо B с ещё одной актуализацией 10:24 письмо C с ещё одной актуализацией 10:30 запускается polling
Список сортируется от новых писем к старым. Первым подходит C. Сервис обрабатывает его и передвигает checkpoint на 10:24. При следующем запуске A и B находятся по другую сторону checkpoint и уже не попадают в обработку.
Важная оговорка: из этого не следует, что итоговые цены обязательно становились неправильными. Если каждый файл являлся полным снимком, C мог полностью заменить A и B, а конечное состояние оставалось корректным.
Но последовательность входящих событий схлопывалась. Для A и B не появлялись собственные результаты обработки, мы не видели, успешно ли они разобрались, чем отличались и какие изменения могли запустить. Система сохраняла последнее состояние, но не полную историю поступивших обновлений.
Polling сохранял последнее состояние, а IMAP IDLE позволил сохранить последовательность событий.
Сейчас этот сценарий кажется очевидным. В тестах старой реализации, однако, были отдельные случаи с одним подходящим письмом и с отсутствием новых писем — но не было пачки из нескольких новых подходящих сообщений внутри одного интервала. Алгоритм делал именно то, что от него формально требовали тесты.
Почему мы пошли в IMAP IDLE
К этому моменту накопилось сразу несколько причин менять модель получения почты.
Письмо могло ждать обработки до 30 минут.
Несколько подходящих писем внутри одного окна могли схлопнуться в самое новое.
Каждый цикл заново подключался к mailbox'ам.
Даже при отсутствии новых сообщений сервис всё равно выполнял пустые проверки.
IMAP IDLE выглядел естественным следующим шагом. Вместо периодического вопроса «не пришло ли что‑нибудь?» клиент оставляет соединение открытым, а сервер сообщает об изменениях в ящике. Новое письмо больше не должно ждать ближайшего получасового окна, а каждое уведомление можно обработать как отдельное событие.
На схеме всё выглядело почти слишком хорошо:

Мы ожидали убрать задержку, перестать делать пустые обходы и восстановить полную последовательность входящих обновлений. Всё это действительно стало достижимым.
Но вместе с 30-минутным таймером исчезла и полезная граница жизненного цикла. Соединение стало долгоживущим. Listener нужно было создавать, останавливать и восстанавливать. Несколько pod теперь могли одновременно услышать один почтовый ящик. Повторная доставка перестала быть исключительным случаем, а корректное завершение приложения перестало быть просто вызовом shutdown().
Именно с этого момента почтовый загрузчик начал становиться распределённой системой.
Первая версия IDLE: listener как динамический ресурс
Первую версию мы строили вокруг простого правила: один физический ящик — один динамический интеграционный flow с IMAP IDLE ресивером.
После полного старта приложения менеджер почтовых конфигов загружал активные, группировал их по email адресу и приводил локальный набор listener'ов к желаемому состоянию. Для нового ящика flow создавался и запускался, для отключённого — останавливался и удалялся.
Изменение настройки во время работы приложения публиковало внутреннее событие. Менеджер снова выполнял синхронизацию, поэтому для включения или выключения почтовой обработки не требовался рестарт всего сервиса.

Долгоживущее соединение неизбежно рвётся: сеть моргает, сервер закрывает сессию, pod переезжает. Поэтому завершение flow с ошибкой не считалось окончательной смертью listener'а. Менеджер планировал его повторный запуск с задержкой, если mailbox всё ещё оставался активным.
На этом этапе IDLE уже решал исходную задачу, но менял гарантии.
Свойство | Polling | Первая версия IMAP IDLE |
|---|---|---|
Задержка | До следующего 30-минутного запуска | Уведомление сервера плюс время обработки |
История событий | Несколько писем могли схлопнуться в самое новое | Каждое уведомление могло запустить отдельную обработку |
Соединение | Короткое, создаётся на каждом проходе | Долгоживущее, требует собственного жизненного цикла |
Восстановление | Следующий запуск по расписанию | Явный restart listener'а и catch‑up после реконнекта |
Последняя строка оказалась самой дорогой. Polling неявно начинал каждую попытку с чистого листа. IDLE сохранял состояние между событиями, а значит, нам пришлось определить, кто этим состоянием владеет и что происходит на каждой границе отказа.
Дешёвый поток — не бесплатная база данных
Письмо, пришедшее через IDLE, быстро покидало сетевой callback и уходило в асинхронную обработку. Для блокирующих операций мы использовали virtual threads. Это удобно: ожидание IMAP, файлового хранилища или базы не требует держать дорогой платформенный поток.
Но virtual threads решают стоимость ожидания, а не пропускную способность системы.
Если одновременно принять много писем и на каждое создать virtual thread, база данных не начнёт выполнять больше запросов в секунду. Пул соединений не расширится сам собой. Парсинг больших файлов не станет дешевле. Мы лишь получим очень эффективный способ поставить больше работы в очередь перед ограниченным ресурсом.
Virtual threads делают блокирующее ожидание дешевле. Они не делают базу данных и downstream‑сервисы быстрее и не ограничивают concurrency.
Поэтому перед тяжёлой обработкой появился отдельный executor с фиксированной границей параллелизма. IDLE мог продолжать принимать события, но количество одновременно работающих конвейеров определялось явно. В тех же изменениях мы добавили более строгую диагностику HikariCP: timeout'ы и обнаружение возможных утечек соединений.
Граница выглядела так:

Это важное различие, модель исполнения и лимит параллелизма — разные решения. Virtual threads отвечают на вопрос «сколько нам стоит ждать?», а bounded executor — «сколько работы downstream вообще разрешено принять одновременно?».
Сразу несколько pod услышали один mailbox
После перехода к IDLE каждый экземпляр приложения на старте видел один и тот же набор активных конфигураций. Без дополнительной координации каждый pod мог решить, что именно он должен создать listener для конкретного ящика.
Для polling это было неприятно, но ограничено отдельными запусками. Для IDLE означало несколько постоянно живущих соединений, которые слушают один ящик и могут одновременно передать одно письмо в обработку.
Пытаться решить это только Kubernetes‑механикой нам не хотелось. Выбор ответственного на весь сервис сделал бы один pod владельцем всех mailbox'ов и создал бы слишком крупную точку отказа и балансировки. Нам требовалось владение на уровне отдельного ящика.
Так появились leases в PostgreSQL.
Для каждого mailbox хранится аренда с идентификатором владельца и сроком действия. Pod пытается атомарно получить свободную или истёкшую аренду. Пока он владелец, он периодически продлевает её и держит локальный listener. При штатной остановке аренда освобождается, при аварийной — перестаёт продлеваться и спустя время становится доступна другому экземпляру.

Lease не означает мгновенный failover: между падением владельца и истечением аренды есть окно. И lease сам по себе не исключает повторную доставку на всех границах отказа. Он отвечает на более узкий вопрос: какой экземпляр сейчас имеет право держать listener для этого mailbox'а.
Так одна таблица в общей базе стала координатором без отдельного coordinator‑сервиса.
Владеть мало — нужно ещё распределять
Первая версия leases обеспечивала эксклюзивность, но не хороший баланс. Представим, что две pod'ы уже разобрали ящики, а затем в кластер добавился третий. Все существующие аренды продолжают исправно продлеваться. Новый экземпляр жив, готов работать — и не получает ничего.
Кроме того, mailbox'ы были неравноценны. Один мог обслуживать одну логическую конфигурацию, другой — множество. Простое деление по количеству ящиков давало честную арифметику, но не обязательно равномерную работу.
Мы добавили heartbeat экземпляров приложения. Каждый pod периодически отмечал в PostgreSQL, что он активен. На основании свежих heartbeat'ов экземпляры понимали текущий размер кластера и могли оценить свою справедливую долю leases.
Затем появился периодический rebalance. Если экземпляр удерживал заметно больше своей доли, он освобождал часть аренд. Свободные mailbox'ы подхватывали менее загруженные pod'ы, включая только что запущенные.
Вес ящика мы приблизили количеством связанных с ним конфигураций. Это не точная модель CPU, размера файлов или частоты писем — только дешёвый proxy, уже доступный в базе. Но он был полезнее, чем считать тяжёлый и пустой mailbox одинаковыми единицами.
Важно, что центрального планировщика у этой схемы нет. Каждый экземпляр видит общее состояние через leases и heartbeat, периодически принимает локальное решение и при необходимости отдаёт часть владения. Поэтому баланс устанавливается не атомарно, а постепенно.
PostgreSQL хранит владение и heartbeat, а pod'ы постепенно сходятся к более равномерному распределению.

Мы получили распределённое владение и отказоустойчивость, не вводя отдельный control plane. Но одновременно открыли следующий класс вопросов: listener можно безопасно передать другой pod'е, только если повторная доставка не превращается в повторное бизнес‑обновление. Значит, после владения нужно было отдельно определить семантику доставки.
Lease не делает обработку однократной
Даже идеально работающий lease не может гарантировать, что письмо попадёт в приложение ровно один раз.
Pod может принять сообщение, начать обработку и завершиться до фиксации результата. Соединение может оборваться между чтением письма и обновлением checkpoint. После reconnect мы снова просматриваем ограниченное окно сообщений, чтобы подобрать то, что могло прийти во время разрыва. При передаче lease новый владелец тоже должен восстановить состояние, не доверяя памяти исчезнувшего процесса.
Во всех этих случаях повторная доставка — нормальный способ восстановления, а не аномалия.
Поэтому мы разделили две гарантии:
lease определяет, кто имеет право слушать mailbox сейчас;
постоянная идемпотентность определяет, нужно ли повторно запускать бизнес‑обработку конкретного письма.
Наша модель — at‑least‑once reception with an idempotent boundary before business processing. Письмо может повторно пересечь входную границу, но перед дорогими и изменяющими состояние действиями мы сверяем его с постоянным журналом обработанных сообщений.
Идентичность письма сложнее одного UID
IMAP даёт сообщению UID, монотонный внутри конкретного mailbox. Для инкрементального catch‑up это намного надёжнее, чем сравнивать только время отправки. timestamp может иметь грубую точность, отличаться от времени получения и зависеть от часов отправителя.
Но UID нельзя вынести из контекста. Сервер публикует UIDVALIDITY — версию пространства UID для folder. Если folder был пересоздан и UIDVALIDITY изменился, прежний UID уже не обязан обозначать то же сообщение.
В нашей реализации UID использовался прежде всего для checkpoint и восстановления. Постоянная идемпотентность сначала искала письмо по mailbox, Message-ID и времени отправки. Если корректного Message-ID не было, fallback строился по mailbox, folder, UIDVALIDITY и UID. Так UID не притворялся глобальным идентификатором, а Message-ID не считался безусловно доступным и уникальным сам по себе.
При первом обнаружении сообщения сервис регистрирует попытку обработки в PostgreSQL. Состояние записи позволяет отличить новый приём от завершённого, повторяемого или уже обрабатываемого. Только после успешной регистрации работа проходит дальше, а терминальный результат и checkpoint фиксируются отдельно.
Упрощённо это выглядит так:

Идемпотентность не превращает распределённую транзакцию в магию. Если внешние эффекты не поддерживают идемпотентный ключ или не входят в одну транзакцию, остаются окна частичного выполнения. Но постоянная запись даёт место, где эти окна можно увидеть, классифицировать и безопасно повторить вместо слепого запуска всего pipeline.
Соблазн второго соединения
К этому моменту у нас уже были leases, идемпотентность и ограничение параллельной обработки. На тяжёлом ящике могло быть много логических конфигураций и заметно больше входящих сообщений, чем на остальных. Естественная гипотеза звучала так: если один IDLE флоу становится точкой последовательного чтения, давайте разрешим выбранному mailbox несколько соединений.
Эксперимент специально ограничили. Multiconnection включался фича‑флагом, применялся только к белому списку выбранных ящиков и имел верхнюю границу числа соединений. Для тяжёлого mailbox существовал и собственный лимит одновременных задач, чтобы он не занял все слоты глобального executor.
На бумаге конструкция казалась безопасной:
lease не даёт разнести mailbox по разным pod'ам + idempotency отсекает повторную доставку + bounded executor защищает downstream = можно параллельно читать один mailbox
В этой формуле не хватало одного участника: IMAP Folder — это stateful ресурс со своим lifecycle.
MimeMessage не обязательно является готовым DTO
Объект MimeMessage, который отдаёт почтовая библиотека, может лениво дочитывать заголовки, multipart‑структуру или содержимое вложения через связанный store и folder. Сам факт, что Java‑объект уже оказался у нас в руках, ещё не означает, что всё письмо загружено в независимую память.
Отсюда возникает неприятная межпоточная гонка:

С несколькими флоу на один mailbox вариантов такого взаимного влияния становится больше. Соединения входят и выходят из IDLE, папки открываются и закрываются, сообщения читаются асинхронно, а retry может пересечься с ещё живой попыткой.
Это не означает, что несколько IMAP‑соединений всегда некорректны. Но в нашей архитектуре потенциальный выигрыш находился на участке чтения почты, а основная тяжёлая работа выполнялась позже — в разборе файлов, базе и downstream‑сервисах. Мы усложнили самый stateful участок системы ради параллелизма до настоящего бутылочного горлышка.
Материализовать, а потом отпускать
Безопасная граница оказалась другой. Пока folder гарантировано открыт, нужно извлечь из MimeMessage всё, что понадобится при последующей обработке и создать независимое представление сообщения. Только такой материализованный объект можно передавать в другой поток, не сохраняя скрытую зависимость от IMAP session.
Вокруг этой границы постепенно сложилось несколько механизмов:
контроль приёма ограничивает число уже принятых, но ещё не завершённых задач;
permit освобождается на всех терминальных путях, а не только при успехе;
in‑flight registry позволяет видеть реально выполняющуюся работу;
checkpoint хранит прогресс чтения ящика и помогает наверстать раззыв после переподключения;
статус повторной попытки отделяет восстанавливаемый сбой от окончательного результата;
graceful shutdown сначала прекращает приём новой работы, а затем предоставляет выполняющимся задачам время на завершение;
интеграция со Spring lifecycle задаёт порядок остановки listener'ов и executor.
Часть этих механизмов потребовалась бы и при одном соединении. Они решают самостоятельные задачи: backpressure, повторную доставку, восстановление и корректное завершение pod. Поэтому после эксперимента мы не стали откатывать весь накопленный слой надёжности.
Для FolderClosedException появился отдельный путь восстановления: попытка не помечалась безусловно завершённой, а могла быть возвращена в состояние, допускающее повтор. Затем multiconnection выключили в конфигурации. После проверки решение закрепили уже архитектурно: удалили feature flag, allowlist, выбор shard'ов, создание нескольких flow и специальные настройки тяжёлых ящиков.
Что осталось в финальной схеме
Финальный вариант снова держит одно IMAP IDLE‑соединение на ответственный mailbox. Но это не возвращение к первой версии IDLE.
Вокруг одного listener'а остались PostgreSQL leases, heartbeat и ребалансировка. На входе остались постоянная идемпотентность, recovery search и checkpoint инфра. Между IMAP и бизнес‑кодом — материализация сообщения, admission control, bounded executor и in‑flight registry. На выходе — явные статусы обработки, а при остановке — управляемый shutdown.

Мы не отказались от параллелизма. Мы перенесли его за точку, где сообщение ещё зависело от состояния IMAP folder.
Безопасная граница параллелизма начинается после материализации сообщения.
Что мы в итоге получили
Главный пользовательский результат прост: новое письмо больше не ждёт очередного запуска scheduler. IDLE позволяет начать приём вскоре после уведомления почтового сервера, а отдельная регистрация каждого сообщения сохраняет историю входящих обновлений вместо схлопывания пачки в самый новый файл.
На уровне кластера у каждого mailbox есть один текущий владелец. При исчезновении pod аренда перестаёт продлеваться, и после её истечения listener может поднять другой экземпляр. Heartbeat и rebalance помогают распределять владение после изменения размера кластера, причём учитывается не только число ящиков, но и приблизительный вес связанных конфигураций.
На уровне обработки всплеска входящих писем больше не означает неограниченный всплеск запросов к базе и downstream‑сервисам. Admission control и bounded executor задают явную емкость, а постоянный журнал идемпотентности защищает бизнес‑границу при reconnect, recovery search и retry.
Наконец, состояние координации стало наблюдаемым. Leases, heartbeat, статусы попыток и in‑flight работа существуют не только в памяти конкретного pod. Их можно диагностировать отдельно: кто слушает mailbox, когда продлевалась аренда, завершилось ли письмо, зависла ли попытка и может ли она быть повторена. Graceful shutdown использует то же явное состояние, чтобы остановить приём и дать уже принятой работе корректно завершиться.
У нас нет сохранённых измерений, из которых можно честно вывести красивое «стало быстрее в N раз» или «надёжность выросла на X%». Поэтому главный результат мы формулируем качественно: уменьшили ожидаемую задержку входа, перестали намеренно пропускать промежуточные события, ограничили downstream concurrency и сделали восстановление явной частью архитектуры.
Чего эта схема не гарантирует
У финального решения остаются важные ограничения.
Это не exactly‑once pipeline. Повторный приём возможен, а идемпотентная граница уменьшает его последствия перед бизнес‑обработкой.
Failover не мгновенный. Новый владелец ждёт, пока прежняя аренда перестанет считаться действующей.
Ребалансировка в конечном итоге будет последовательной: после scale‑up или изменения нагрузки распределение сходится постепенно.
UID search и checkpoint инфра присутствуют в коде, но конкретный режим поиска зависит от конфигурации. В текущем prod‑конфиге репозитория UID search выключен, это не следует скрывать за общей схемой.
Детали IMAP IDLE, reconnect и server‑side поиск зависят от реализации почтового сервера. Поведение нужно проверять на том сервере, с которым система действительно работает.
Идемпотентная запись не делает атомарными произвольные внешние эффекты. Для каждого downstream всё равно нужен ответ на вопрос, что произойдёт при повторе после частичного выполнения.
Эти ограничения не обесценивают архитектуру. Они задают её настоящий контракт — тот, на который можно опираться при следующем сбое.
Вывод

Мы начинали с простой задачи: перестать проверять почту раз в 30 минут. Потом появились долгоживущие соединения, leases, heartbeat, идемпотентность, backpressure, checkpoints и попытка распараллелить ещё немного.
А в итоге самым полезным изменением оказалось не то, которое добавило больше параллелизма, а то, которое помогло понять, где этот параллелизм вообще должен начинаться.

