Под получил SIGTERM. HTTP-запросы Spring сольёт сам, Kafka-контейнеры остановит сам, а фоновая задача в эту секунду может стоять посреди чужого вызова: деньги у платёжного провайдера уже списаны, строки о списании в базе ещё нет. Откат транзакции тут не поможет — откатывать нечего, повтор при следующем запуске не поможет — повторять некому. Фоновые задачи — единственное место в приложении, где остановка теряет результат молча: без ошибки в логе и без 5xx клиенту.
Таких задач в Spring три вида — @Scheduled, @Async и планировщик, который гоняет outbox-relay, — и останавливаются они по-разному: одних ждут, других прерывают, третьих просто перестают ждать. Восстановить недоделанное получается только у одного из трёх.
Разница не в длине ожидания (30, 20 и 25 с) и даже не в том, прервут задачу или просто перестанут её ждать, — а в том, где лежит незаконченная работа: в памяти пула её после остановки не восстановить, в строке outbox её заберёт SKIP LOCKED при следующем запуске.
Как Spring останавливает @Scheduled-задачи
Relay отправляет пакет из пятидесяти событий, на двадцатом приходит SIGTERM. Если ничего не настраивать, планировщик поведёт себя аккуратнее, чем принято думать. На событие закрытия контекста он перестаёт ставить в расписание новые итерации, а текущей даёт доработать: фаза остановки не закончится, пока в пуле есть работающая задача. Предел — общий для всей группы spring.lifecycle.timeout-per-shutdown-phase, 30 секунд по умолчанию. Что не успело за это время, получит прерывание позже, когда очередь дойдёт до уничтожения бинов: пул закрывается через shutdownNow(), а тот шлёт потоку interrupt().
Тридцати секунд хватает почти всегда: итерация relay — это пятьдесят отправок в Kafka по несколько миллисекунд каждая. Не хватает, когда внутри итерации завис внешний вызов. Тогда ожидание продлевают отдельной настройкой:
spring:
task:
scheduling:
shutdown:
await-termination: true
await-termination-period: 25s
pool:
size: 4
Только меняет она не «ждать или не ждать», а где именно приложение ждёт, — и разница неочевидная. Планировщик с этим флагом на событие закрытия контекста не реагирует вовсе и продолжает работать по расписанию; свою фазу остановки он проходит мгновенно. Ожидание переехало в самый конец, на уничтожение бинов: там новые задачи запрещаются, а текущей отводятся те самые 25 секунд. То есть итерацию планировщик дожидается уже после того, как Spring слил HTTP-запросы, а не до. Для бюджета остановки это значит, что 25 секунд складываются с временем дрейна HTTP, а не идут параллельно ему.
Не уложилась итерация и в 25 секунд — прерывания не будет. Spring запишет в лог Timed out while waiting for executor 'taskScheduler' to terminate и пойдёт закрывать контекст дальше, а задача останется работать в непрерванном потоке, пока её не оборвёт закрытый пул соединений или остановка JVM по окончании хуков. await-termination-period — не гарантия, что задачу доведут до конца. Это обещание подождать, и не больше.
Отдельная ловушка — потоки планировщика и пулов Spring не демоны. При SIGTERM это ничего не меняет: JVM дожидается хуков остановки и гасит процесс вместе со всеми потоками. Но если приложению положено завершиться самому — задание Kubernetes Job, консольная утилита на Spring Boot, — один-единственный @Scheduled держит JVM живой навечно: main вернулся, а поток scheduling-1 не демон и ждёт следующей итерации. Job висит в Running до activeDeadlineSeconds или до ручного удаления. Лечится одной строкой — System.exit(SpringApplication.exit(context)): закрытие контекста останавливает пул, а exit гасит JVM независимо от того, кто ещё жив.
Три пода — три прогона
@Scheduled ничего не знает о соседях. Три реплики — три параллельных прогона одной и той же задачи с точностью до расхождения часов. Outbox-relay это переживает: FOR UPDATE SKIP LOCKED делит строки между подами, и каждый забирает свои. Обычная задача — «раз в ночь пересчитать остатки», «раз в час выслать отчёт» — выполнится трижды, и отчёт уйдёт трижды.
Для таких задач ставят распределённую блокировку, чаще всего ShedLock: строку в таблице shedlock с полями lock_until и locked_by, которую задача захватывает перед стартом. Нужны @EnableSchedulerLock и бин LockProvider над той же базой (JdbcTemplateLockProvider), отдельной инфраструктуры не требуется.
@Scheduled(cron = "0 0 3 * * *")
@SchedulerLock(name = "recalculateStock", lockAtMostFor = "PT20M", lockAtLeastFor = "PT1M")
public void recalculateStock() {
stockService.recalculate();
}
lockAtMostFor — это и есть связь с остановкой. Под, убитый посреди пересчёта, блокировку не снимет: она отпустится сама, когда истечёт срок. Поставили PT20M при пересчёте на пять минут — после неудачного выката задача не запустится ни на одной реплике ещё пятнадцать минут, и это нормальная цена. Поставили PT2M при пересчёте на пять — вторая реплика начнёт пересчёт поверх первой на третьей минуте, пока первая ещё работает. Срок выбирают заведомо больше самого долгого прогона, а не «примерно как обычно». lockAtLeastFor защищает от обратного: часы на подах разошлись на секунду, первый закончил за полсекунды и отпустил блокировку, второй по своим часам ещё не стартовал — и стартует.
@Async: где теряются деньги
Метод chargeCustomer помечен @Async, пул под него никто не настраивал. Вызов к платёжному провайдеру прошёл — деньги списаны. До сохранения результата в базу поток получает прерывание: приложение останавливается. В базе нет записи о платеже, деньги ушли. Идемпотентность тут не спасает: она защищает от повторных попыток, а повторять некому — результат жил только в памяти пула, и с остановкой его больше нет.
Без своего бина Spring Boot 3 подставляет под @Async автонастроенный applicationTaskExecutor: ThreadPoolTaskExecutor на восемь основных потоков с неограниченной очередью и именами task-1, task-2. Восемь потоков — это и потолок: очередь без предела, поэтому до max-size дело не доходит никогда, лишние задачи копятся в очереди. При остановке пул ведёт себя как планировщик по умолчанию: получив событие закрытия, перестаёт брать новое, запущенное дожидается в фазе остановки в общие 30 секунд, а на уничтожении бинов зовёт shutdownNow() — запущенным прилетает interrupt(), очередь невыполненных выбрасывается целиком. Для chargeCustomer это означает: успел за 30 секунд — хорошо; завис на провайдере — прерывание, и результат пропал.
Чтобы отвести пулу собственный срок, объявляют бин исполнителя явно:
@Bean
public ThreadPoolTaskExecutor asyncExecutor() {
var executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(8);
executor.setMaxPoolSize(16);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("async-");
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(20);
return executor;
}
Два флага, и нужны оба. setWaitForTasksToCompleteOnShutdown(true) меняет shutdownNow() на shutdown(): запущенные задачи не прерывают, очередь дорабатывают. Но сам по себе он ничего не ждёт — без setAwaitTerminationSeconds пул закроется мгновенно, и задачи останутся работать в потоках, которых уже никто не дожидается. setAwaitTerminationSeconds(20) — сколько приложение готово стоять на уничтожении бинов; как и у планировщика, это ожидание идёт после дрейна HTTP и по истечении срока просто кончается, без прерывания. Для автонастроенного пула те же два параметра зовутся spring.task.execution.shutdown.await-termination и await-termination-period. Ни initialize(), ни destroyMethod писать не нужно: ThreadPoolTaskExecutor сам реализует интерфейсы запуска и уничтожения. Ручной initialize() в @Bean-методе даже вреден — Spring вызовет инициализацию ещё раз и создаст второй пул поверх первого.
С виртуальными потоками правила меняются. spring.threads.virtual.enabled: true подменяет applicationTaskExecutor на SimpleAsyncTaskExecutor, и пула там нет вовсе: на каждую задачу — новый виртуальный поток. Без await-termination такой исполнитель при остановке не делает ничего: не ждёт и не прерывает, а виртуальные потоки всегда демоны и умирают вместе с JVM, где бы ни стояли. С await-termination: true порядок обратный привычному: close() сначала шлёт interrupt() всем активным задачам и только потом ждёт до await-termination-period. На виртуальных потоках «подождать» означает «прервать и подождать, пока прерванные доработают», а не «дать доработать, не трогая». Планировщик при этом флаге тоже становится виртуальным (SimpleAsyncTaskScheduler) и закрывается так же: отменяет запланированное, прерывает текущее, ждёт.
Если прерывание до задачи всё-таки доходит — пул по умолчанию, виртуальные потоки, — блокирующие вызовы вроде Future.get() и Thread.sleep отвечают на него InterruptedException, и глотать его нельзя:
@Component
@RequiredArgsConstructor
public class AsyncEmailSender {
private final EmailClient emailClient;
@Async
public CompletableFuture<Void> sendAsync(String to, String subject, String body) {
try {
Future<Void> sent = emailClient.send(to, subject, body);
sent.get(10, TimeUnit.SECONDS);
return CompletableFuture.completedFuture(null);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return CompletableFuture.failedFuture(e);
} catch (Exception e) {
return CompletableFuture.failedFuture(e);
}
}
}
Thread.currentThread().interrupt() в catch возвращает потоку флаг прерывания: get() его сбросил, и без этой строки следующий блокирующий вызов в той же задаче повис бы как ни в чём не бывало, а пул так и не узнал бы, что поток просили остановиться.
Когда @Async не нужен
@Async уместен там, где результат можно потерять: прогреть кэш, отправить необязательное уведомление, записать метрику. Как только работа заканчивается строкой в базе — списание, заказ, письмо, которое обязано уйти, — фоновый поток не нужен вовсе. Намерение кладут в таблицу той же транзакцией, что и бизнес-данные, а выполняет его воркер по расписанию. Тогда прерывание значит только «продолжим в следующем запуске», а не потерю результата. Это и есть outbox.
Outbox-relay: цикл, который переживает остановку
Заказ создан — событие об этом должно уйти в Kafka. Отправлять прямо из транзакции запроса нельзя: транзакция ещё может откатиться, а событие уже улетело. Отправлять после коммита — можно упасть между коммитом и отправкой, и склад о заказе не узнает никогда. Outbox разводит это на два шага: запрос пишет событие в таблицу outbox_event той же транзакцией, что и заказ, а отдельный relay по расписанию забирает строки без отметки, шлёт их и отметку ставит.
Две транзакции и одна таблица между ними. Запрос кладёт событие рядом с заказом и ни о какой Kafka не знает; relay забирает пакет под блокировкой и ставит отметки все разом, на коммите. Обрыв на сороковом send() откатывает отметки, но не отправки: строки вернутся в выборку, сорок событий уедут второй раз.
@Component
@RequiredArgsConstructor
public class OutboxRelay {
private final DSLContext dsl;
private final KafkaTemplate<String, Object> kafkaTemplate;
@Scheduled(fixedDelay = 500)
@Transactional
public void publish() {
var batch = dsl.selectFrom(OUTBOX_EVENT)
.where(OUTBOX_EVENT.PUBLISHED_AT.isNull())
.orderBy(OUTBOX_EVENT.ID)
.limit(50)
.forUpdate().skipLocked()
.fetch();
for (var row : batch) {
kafkaTemplate.send(row.getTopic(), row.getPartitionKey(), row.getPayload()).join();
row.setPublishedAt(OffsetDateTime.now());
row.store();
}
}
}
FOR UPDATE SKIP LOCKED делает две вещи сразу. FOR UPDATE блокирует выбранные строки до конца транзакции relay, SKIP LOCKED велит соседнему relay в другом поде не ждать этих строк, а взять следующие. Так пакеты делятся между репликами без координатора, и тот же приём защищает от остановки. При SIGTERM посреди publish() планировщик дожидается вызова; успела транзакция зафиксироваться — строки помечены. Не успела, оборвалась вместе с пулом соединений — транзакция откатилась, блокировка снята, published_at пуст. Следующий запуск или соседний под захватит эти же строки и отправит их снова.
Мелочь, на которой спотыкаются: тип у setPublishedAt. Колонка published_at объявлена как timestamptz, и jOOQ генерирует под неё сеттер на OffsetDateTime, а не на Instant; для колонки без зоны (timestamp) это был бы LocalDateTime.
Потеряться событие не может: пока строка без отметки, её заберут снова. А вот продублироваться — запросто, и это цена схемы. В Kafka события уходят по одному, а отметки ложатся в базу все разом, на коммите. Откатилась транзакция после того, как первые сорок улетели, — следующий запуск отправит эти сорок повторно. Поэтому outbox всегда идёт в паре с защитой от повторов на стороне получателя: об этом — идемпотентность при остановке.
Ядовитое событие
Тот же цикл, что спасает при остановке, умеет вставать намертво. Событие с id 57 весит полтора мегабайта — кто-то положил в payload весь заказ с картинками, — а продюсер не пропускает записи больше max.request.size, по умолчанию 1 МБ. send() падает с RecordTooLargeException, транзакция откатывается вместе с отметками первых шестнадцати, и через 500 мс relay берёт тот же пакет: ORDER BY id вернёт те же пятьдесят строк, начиная с тех же шестнадцати. Отправка событий остановилась на всех подах, хотя все они живы, а в логе одна и та же ошибка раз в полсекунды.
Лечат счётчиком попыток и отсечкой: в таблицу добавляют attempts и last_error, в выборку — .and(OUTBOX_EVENT.ATTEMPTS.lt(10)), а неудачу записывают в той же транзакции и пакет заканчивают:
for (var row : batch) {
try {
kafkaTemplate.send(row.getTopic(), row.getPartitionKey(), row.getPayload())
.get(5, TimeUnit.SECONDS);
row.setPublishedAt(OffsetDateTime.now());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
} catch (ExecutionException | TimeoutException e) {
row.setAttempts(row.getAttempts() + 1);
row.setLastError(e.getMessage());
row.store();
break;
}
row.store();
}
break вместо исключения — намеренно: метод завершается штатно, транзакция фиксируется, и с ней сохраняются отметки уже отправленных событий и счётчик у виновника. Десять неудач по 500 мс — и через пять секунд событие 57 выпадает из выборки, очередь идёт дальше. Само событие никуда не делось: строку с attempts = 10 показывает отдельная метрика, и разбирает её человек. Порядок событий по ключу заказа 57 при этом сломан — следующие по нему уже уехали, — так что «положить обратно» после починки payload недостаточно, получателю понадобится сверка.
Пакет, который не успевает
Пятьдесят отправок по 5–10 мс — четверть секунды, влезает в любое окно. Пакет перестаёт влезать не от размера, а от одного зависшего send(): у продюсера delivery.timeout.ms по умолчанию 120 секунд — столько он перепосылает запись, пока брокер не отвечает, и ровно столько провисит join(). Одно событие на недоступном брокере съедает весь бюджет остановки, а после него в пакете ещё сорок девять. Поэтому в relay ждут результат с собственным пределом — те самые get(5, TimeUnit.SECONDS) в листинге выше: пять секунд на событие, дальше это неудача отправки, попытка засчитана, пакет закончен.
Второй рычаг — признак остановки между событиями. Планировщик с await-termination на событие закрытия контекста не реагирует, но приложение — может: слушатель ContextClosedEvent выставляет volatile boolean stopping, а цикл проверяет его перед каждым send() и выходит через break. Транзакция фиксируется, отметки отправленных сохраняются, дублей от этого пакета не будет. Уменьшать сам пакет имеет смысл только после этих двух правок: пакет в десять строк с зависшим на две минуты join() ничем не лучше пакета в пятьдесят.
Когда outbox не берут
Outbox стоит денег, которых не видно в примере на пятьдесят строк. Таблица растёт на каждое событие, и чистить её — отдельная задача (DELETE … WHERE published_at < now() - interval '7 days', тоже по расписанию). Relay — ещё один компонент со своими метриками и своим лагом: минимум 500 мс от коммита до Kafka. Получатель обязан уметь в повторы. Это оправдано, когда событие терять нельзя: заказ, оплата, смена статуса. Когда потерю переживут — «товар посмотрели» для аналитики, сброс кэша с TTL, метрика, — событие шлют прямо из кода после коммита, и остановка посреди отправки означает одну потерянную запись в счётчике, а не расхождение денег.
Частая ошибка: бесконечный цикл внутри @Scheduled
@Scheduled(fixedDelay = 100)
public void relay() {
while (true) {
publish();
}
}
Пауза между итерациями — единственное место, где планировщик может сказать «хватит». Цикл внутри итерации эту паузу съедает: планировщик честно ждёт конца итерации, которого не будет.
Планировщик умеет останавливать задачу только между итерациями: не запускает следующую и ждёт, пока кончится текущая. Итерация с while (true) не кончается никогда, и ждать её планировщик будет весь отведённый срок. Прерывание, которое прилетит на уничтожении бинов, поможет, только если внутри цикла есть блокирующий вызов, который на него отвечает, — Thread.sleep, Future.get, join() на результате send(). Цикл, который крутит запросы к базе, прерывания не заметит: JDBC-драйвер PostgreSQL на interrupt() не реагирует, и поток доживёт до остановки JVM. Правильно — одна итерация обрабатывает один пакет и завершается, а частоту задаёт fixedDelay: следующая итерация стартует через 500 мс после конца предыдущей, и пустой прогон при пустом outbox стоит один быстрый SELECT.
Глубже: что изменилось в Spring 6.1: исполнители как SmartLifecycleрасширенное
Поведение планировщика и @Async выше описано для текущих версий, Spring Boot 3.2 и новее с Framework 6.1. Старые статьи и ответы на форумах описывают модель Boot 2.x, и расхождение стоит назвать явно, чтобы не искать в своём приложении того, чего в нём больше нет.
До 6.1 ThreadPoolTaskExecutor и ThreadPoolTaskScheduler в остановке контекста не участвовали: они узнавали о ней только при уничтожении бинов, в самом конце, и до этого момента принимали и запускали задачи как ни в чём не бывало. Ожидание waitForTasksToCompleteOnShutdown тоже происходило при уничтожении. Отсюда классическая картина: веб-сервер уже остановлен, а планировщик ставит новые итерации ещё десять секунд.
С 6.1 исполнители и планировщик реализуют SmartLifecycle с фазой Integer.MAX_VALUE / 2. Остановка идёт по фазам от большей к меньшей, веб-сервер останавливается в фазе выше, поэтому порядок такой: веб-сервер перестал принимать запросы и дренирует их, затем исполнители получают ранний сигнал остановки и перестают принимать новые задачи, и только потом, при уничтожении бинов, идёт ожидание уже запущенных, если оно включено. Планировщик с этого момента не ставит новых итераций, что и описано в разделе выше. Отсюда два новых свойства. acceptTasksAfterContextClose (в Boot spring.task.execution.shutdown.accept-tasks-after-context-close) разрешает исполнителю принимать задачи и после раннего сигнала: это нужно, когда запрос, который дренируется, отправляет @Async-задачу, и без свойства она получит TaskRejectedException. И strictEarlyShutdown, который наоборот прерывает запущенные задачи уже на раннем сигнале, а не при уничтожении; по умолчанию выключен.
Что из этого следует для чтения статей раздела: spring.lifecycle.timeout-per-shutdown-phase ограничивает каждую фазу, а ожидание задач исполнителя живёт в своём await-termination-period при уничтожении, и эти два числа складываются в бюджет по-разному, чем в Boot 2.x. Своя обёртка над пулом, написанная под старую модель, с флагом завершения и ручным ожиданием, в 6.1 дублирует то, что фреймворк уже делает, и мешает ему; такие обёртки убирают, о чём раздел про собственный флаг завершения в статье про конфигурацию.
Коротко
- Планировщик по умолчанию дожидается текущей итерации в фазе остановки (до 30 с) и прерывает её только при уничтожении бинов;
await-terminationпереносит ожидание в самый конец, после дрейна HTTP, и по истечении срока просто перестаёт ждать. @Asyncбез своего бина живёт вapplicationTaskExecutorна восемь потоков и получаетshutdownNow(); чтобы задачу не прервали, нужны оба флага —setWaitForTasksToCompleteOnShutdown(true)иsetAwaitTerminationSeconds. На виртуальных потоках «подождать» значит «прервать и подождать».- Ждать — не значит сохранить: результат, который жил в памяти пула, после остановки не восстановить. Работа, которая кончается строкой в базе, идёт через outbox, а не через
@Async. - Outbox переживает обрыв за счёт
FOR UPDATE SKIP LOCKED: откат снимает блокировку, следующий запуск забирает строки снова. Цена — повторы, которые отсеивает получатель, счётчик попыток против ядовитого события и свой таймаут наsend()вместоjoin(). - Три реплики — три прогона
@Scheduled; outbox делит строки сам, остальным нужен ShedLock сlockAtMostForдлиннее самого долгого прогона. while (true)внутри@Scheduledпланировщик остановить не может: частоту задаётfixedDelay, а итерация всегда завершается сама.- С Spring 6.1 исполнители и планировщик это
SmartLifecycleс фазойMAX_VALUE / 2: перестают принимать задачи после дренажа веб-сервера, ожидание идёт при уничтожении;accept-tasks-after-context-closeдля задач из дренируемых запросов, обёртки под модель Boot 2.x убирают.
Что почитать дальше
- JVM и Spring: базовая конфигурация graceful shutdown — откуда берутся фаза остановки,
timeout-per-shutdown-phaseи порядок уничтожения бинов, на которые здесь всё опирается. - Идемпотентность при остановке — как получатель отсеивает те сорок событий, которые relay отправит дважды.
- Бюджеты и наблюдаемость — как уложить 25 секунд планировщика и 20 секунд пула в общие 60 и увидеть в метриках, что не уложились.
- Outbox и Inbox — сам паттерн целиком: таблица, relay, inbox на стороне получателя и порядок событий по ключу.