← назад к разделу

Под получил SIGTERM. HTTP-запросы Spring сольёт сам, Kafka-контейнеры остановит сам, а фоновая задача в эту секунду может стоять посреди чужого вызова: деньги у платёжного провайдера уже списаны, строки о списании в базе ещё нет. Откат транзакции тут не поможет — откатывать нечего, повтор при следующем запуске не поможет — повторять некому. Фоновые задачи — единственное место в приложении, где остановка теряет результат молча: без ошибки в логе и без 5xx клиенту.

Таких задач в Spring три вида — @Scheduled, @Async и планировщик, который гоняет outbox-relay, — и останавливаются они по-разному: одних ждут, других прерывают, третьих просто перестают ждать. Восстановить недоделанное получается только у одного из трёх.

SIGTERM в 0 с · окно ожидания фоновых задач из бюджета 60 с0 с10 с20 с25 с30 с @Async без настроексрок общий: timeout-per-shutdown-phase ждём завершения до 30 с не успел — shutdownNow() прерывает поток: платёж ушёл, записи в базе нет @Async · awaitTerminationSeconds(20)waitForTasksToCompleteOnShutdown(true)ждём завершения до 20 с перестали ждатьне успел — результат в памяти, при старте его не подхватит никто outbox-relay · await-termination 25 с@Scheduled(fixedDelay = 500) · пакет 50ждём завершения до 25 с перестали ждатьстроки без published_at и без блокировки — SKIP LOCKED возьмёт их снова Ожидание покупает время, а не сохранность результатаПрерывание безопасно там, где незаконченная работа лежит строкой в базе

Разница не в длине ожидания (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 по расписанию забирает строки без отметки, шлёт их и отметку ставит.

Строка outbox рождается вместе с заказом и живёт до отметки published_at HTTP-запрос · одна транзакция INSERT orders INSERT outbox_event COMMIT вместе или никак outbox_event id 41…90 · published_at IS NULL · 50 строк relay · fixedDelay 500 мс · транзакция на пакет relay SELECT … SKIP LOCKED send() ×50 published_at=now() COMMIT SIGTERM: обрыв после 40-го send() ROLLBACK: 50 строк снова без отметки и без блокировки следующий запуск или соседний под забирает их снова — 40 событий уйдут дважды Потерять событие нельзя, продублировать — легко поэтому у получателя стоит защита от повторов

Две транзакции и одна таблица между ними. Запрос кладёт событие рядом с заказом и ни о какой 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();
    }
}
Планировщик останавливает только между итерациями fixedDelay итерация 500 мс итерация 500 мс итерация стоп новую итерацию не запустят, текущую дождутся while (true) publish() → publish() → publish() → … выхода нет вклиниться некуда — ждёт весь срок; interrupt увидит лишь sleep/get/join

Пауза между итерациями — единственное место, где планировщик может сказать «хватит». Цикл внутри итерации эту паузу съедает: планировщик честно ждёт конца итерации, которого не будет.

Планировщик умеет останавливать задачу только между итерациями: не запускает следующую и ждёт, пока кончится текущая. Итерация с 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 убирают.

Что почитать дальше