При остановке приложения Kafka-потребитель и продюсер нужно завершать аккуратно. Иначе одна из двух бед: либо сообщения обрабатываются повторно при следующем запуске, либо часть сообщений, которые уже отправил продюсер, так и не добирается до брокера.
Разберём оба случая по порядку.
Главный риск — на стороне потребителя, поэтому начнём с него: посмотрим на одной шкале, сколько времени занимает listener и сколько на него отведено.
Окно одно и то же — 20 секунд из timeout-per-shutdown-phase. Три HTTP-вызова с повторами тянут 65 секунд и обрываются на двадцатой: offset не зафиксирован, а платёж уже ушёл. Запись в базу и строка в outbox укладываются в миллисекунды, и остановка проходит чисто.
Что происходит с потребителем при остановке
Kafka-потребитель читает сообщения пачками. После обработки он фиксирует прогресс в брокере — это называется «закоммитить offset». Если приложение убить до коммита, при следующем запуске потребитель прочитает те же сообщения заново (replay).
Spring управляет потребителями через ConcurrentMessageListenerContainer. При закрытии контекста он просит каждый контейнер остановиться: тот дорабатывает текущую пачку, коммитит offset и закрывает KafkaConsumer. Контейнеры Kafka останавливаются в числе первых — раньше, чем начинается слив HTTP-запросов: у контейнеров фаза остановки Integer.MAX_VALUE - 100, у слива веб-сервера Integer.MAX_VALUE - 1024, а чем больше число, тем раньше остановка. Поэтому окно отсчитывается практически от самого SIGTERM.
Ключевая настройка — время ожидания:
spring:
lifecycle:
timeout-per-shutdown-phase: 20s
kafka:
listener:
ack-mode: MANUAL
spring.lifecycle.timeout-per-shutdown-phase (по умолчанию 30 с) — максимум, сколько Spring ждёт остановки группы бинов, в которую входят контейнеры Kafka. Отдельного свойства spring.kafka.listener.shutdown-timeout в Spring Boot нет: у контейнера есть своя настройка ContainerProperties.shutdownTimeout (10 с), но она срабатывает только при ручном вызове container.stop() — при закрытии контекста границу задаёт timeout-per-shutdown-phase.
Важная оговорка про эти 20 секунд: свойство одно на всё приложение. Поставив его ради Kafka, вы тем же движением урезаете и окно дрейна HTTP, и ожидание планировщика — а групп в приложении несколько, и гасятся они одна за другой. Собирая конфигурацию из разных статей этого раздела, выбирайте одно значение на всех; 30 секунд по умолчанию — разумная отправная точка, 20 здесь взяты только чтобы окно было видно на схеме.
Если listener не уложился в это окно, Spring его не прерывает: он пишет в лог, что группа не остановилась, и продолжает закрывать контекст. Дальше процесс завершается, поток обработки умирает на середине — offset не зафиксирован, и при следующем старте сообщения придут снова.
Почему нельзя делать долгие операции внутри listener
Listener обрабатывает сообщение и укладывается в окно остановки. Проблема возникает, когда внутри listener идут несколько HTTP-вызовов с повторными попытками:
// Опасно — суммарное время 65s, в окно остановки 20s не влезет
@KafkaListener(...)
public void onConfirmed(OrderConfirmedEvent event, Acknowledgment ack) {
paymentClient.charge(...); // 5s + повторы 30s
notificationClient.send(...); // 3s + повторы 15s
analyticsClient.track(...); // 2s + повторы 10s
ack.acknowledge();
}
Ту же длительность ограничивают две настройки самого потребителя, и знать их надо обе, потому что одна касается остановки, а вторая — обычной работы.
max.poll.records (по умолчанию 500) — сколько сообщений отдаётся за один запрос к брокеру. Это и есть главный рычаг остановки: контейнер, получив просьбу завершиться, дообрабатывает текущую пачку и только потом останавливается. Пятьсот сообщений по сто миллисекунд — это пятьдесят секунд, которых в бюджете нет. Пятьдесят сообщений — пять секунд, и всё укладывается.
spring:
kafka:
consumer:
max-poll-records: 50
properties:
max.poll.interval.ms: 120000
max.poll.interval.ms (по умолчанию 5 минут) — сколько времени группа готова ждать между запросами, прежде чем счесть потребителя зависшим и забрать у него партиции. Он про обычную работу: долгая обработка пачки выглядит для группы как зависание, начинается перераспределение, и наступает знакомая картина — «потребление встало, а в журнале строки про присоединение к группе».
Связь между ними простая: max.poll.records умножить на время обработки одного сообщения должно быть заметно меньше max.poll.interval.ms — и одновременно меньше бюджета остановки. Если оба условия не выполняются одновременно, обработку выносят из слушателя: слушатель только принимает сообщение и складывает задачу, а долгую работу делает отдельный исполнитель, у которого свои правила остановки.
В худшем случае listener занимает 65 секунд. Окно в 20 секунд истечёт, Spring перестанет ждать и дозакроет контекст, процесс завершится — offset не зафиксируется. При перезапуске событие придёт снова — и если первая попытка paymentClient.charge уже дошла до платёжного сервиса, получится двойное списание.
Решение — listener делает только быструю локальную работу, а «тяжёлые» операции откладывает через outbox-паттерн:
@KafkaListener(...)
@Transactional
public void onConfirmed(OrderConfirmedEvent event, Acknowledgment ack) {
if (!processedEventRepository.tryMarkProcessed(event.eventId(), "billing")) {
ack.acknowledge();
return;
}
billingService.recordChargeIntent(event.orderId(), event.totalAmount());
outboxRepository.append(
new ChargePaymentRequested(event.orderId(), event.totalAmount())
);
ack.acknowledge();
}
Listener сохраняет намерение в базу и кладёт событие в outbox — всё в одной локальной транзакции, завершается за миллисекунды. Реальный HTTP-вызов к платёжному сервису делает отдельный процессор outbox. Shutdown проходит чисто.
Режим подтверждения сообщений — ack-mode
«Зафиксировать offset» значит сказать брокеру: всё до этой позиции обработано, при перезапуске начинай дальше. Остановился процесс до фиксации — сообщения придут повторно; после фиксации, но до конца обработки — они потеряны. ack-mode выбирает момент фиксации, а значит и то, что случится при остановке.
Одна пачка проходит четыре точки, и ack-mode выбирает, в какой из них offset уедет брокеру, а MANUAL_IMMEDIATE фиксирует раньше всех, прямо в момент acknowledge внутри обработки.
| ack-mode | Когда фиксируется offset | Поведение при shutdown |
|---|---|---|
| RECORD | После каждого сообщения | Точно, но больше запросов к брокеру |
| BATCH | После обработки всей пачки предыдущего poll() | Либо вся пачка закоммичена, либо нет — чисто |
| TIME | Когда истёк ack-time | Остаток коммитится при остановке контейнера |
| COUNT | Когда накопилось ack-count подтверждений | Остаток коммитится при остановке контейнера |
| COUNT_TIME | Что из этих двух наступит раньше | Остаток коммитится при остановке контейнера |
| MANUAL | Вы зовёте ack.acknowledge(), но коммит уходит, когда обработана вся пачка | Подтверждение может не успеть превратиться в коммит |
| MANUAL_IMMEDIATE | В момент ack.acknowledge(), если он вызван в потоке потребителя | Коммит уже у брокера — повтора не будет |
Для большинства случаев подходит BATCH: один коммит на всю пачку, минимум накладных расходов, при остановке поведение предсказуемо — либо вся пачка обработана и зафиксирована, либо придёт на повтор целиком.
На последние две строки стоит посмотреть внимательно, потому что название обманывает. ack.acknowledge() в режиме MANUAL ничего не коммитит: подтверждения складываются в очередь и уезжают к брокеру, только когда обработана вся пачка предыдущего poll(). Для остановки это и есть ключевая разница — вызвали acknowledge(), а контейнер в этот момент перестали ждать, и коммита не случилось: сообщение придёт снова. Коммит ровно в момент вызова даёт только MANUAL_IMMEDIATE.
Примеры в этой статье написаны под MANUAL: аргумент Acknowledgment Spring передаёт в listener только в ручных режимах. При BATCH такой метод упадёт с IllegalStateException: No Acknowledgment available as an argument.
immediate-stop: дообрабатывать пачку или бросить сразу
У контейнера слушателя есть настройка, которая прямо про остановку и о которой почти не знают:
spring:
kafka:
listener:
immediate-stop: false # значение по умолчанию
Со значением false (по умолчанию) контейнер, получив просьбу остановиться, дообрабатывает текущую пачку целиком и только потом закрывает потребителя. Это правильное поведение: пачка либо обработана и зафиксирована, либо придёт на повтор — без разрывов посередине.
Со значением true контейнер останавливается сразу после текущего сообщения, не доводя пачку до конца. Остаток пачки не обрабатывается и не фиксируется — эти сообщения придут заново другому поду.
Когда включают true: когда обработка одного сообщения дорогая и пачка длинная, а бюджет остановки маленький; когда обработчик гарантированно идемпотентен, и повтор ничего не стоит. В остальных случаях значение по умолчанию лучше — потому что «бросить на середине» означает больше повторов, а значит больше нагрузки на идемпотентность.
Отдельно стоит знать, что остановка контейнера — это не одно мгновение: Spring вызывает stop() и ждёт его завершения в пределах фазы остановки (spring.lifecycle.timeout-per-shutdown-phase). Если пачка не уложилась, дальше сработает общий принудительный путь, и незафиксированные сообщения придут повторно.
Самая опасная ошибка — enable.auto.commit
Сразу важная оговорка, которая меняет вес всего раздела: Spring Kafka сам выставляет enable.auto.commit в false, если вы не указали иначе. То есть по умолчанию описанной ниже беды у вас нет — она возможна только если настройку включили руками, обычно перенеся конфигурацию из примера на «голом» клиенте или из другого проекта. Поэтому раздел стоит читать как «почему это не надо включать», а не как «почему у вас всё сломано».
Режима AUTO в ContainerProperties.AckMode нет: автокоммит живёт отдельной настройкой самого потребителя. Есть настройка, которая кажется удобной, но при остановке приложения приводит к потере сообщений:
# Опасно
spring:
kafka:
consumer:
enable-auto-commit: true
auto-commit-interval-ms: 5000
Вот как происходит потеря — и происходит она не так, как обычно рассказывают. Никакого фонового потока, который коммитит offset'ы, у потребителя нет. Коммит делает сам poll(), на том же потоке, если с прошлого раза прошло больше auto.commit.interval.ms. И отмечает он offset'ы тех записей, которые отдал предыдущий poll(), — не спрашивая, довели вы их обработку до конца или нет.
Отсюда и сценарий: следующий poll() случился раньше, чем дообработана предыдущая пачка. Он забрал сообщения 200 и дальше, а заодно отметил прогресс на 200 — хотя реально обработаны только 100–149. Приложение в этот момент останавливают. После перезапуска потребитель начнёт читать с 200, и сообщения 150–199 потеряны навсегда.
Правильно — enable-auto-commit: false и управлять коммитами явно через ack-mode: BATCH или MANUAL. В Spring Kafka это уже так по умолчанию: достаточно не включать настройку.
Продюсер — flush при остановке
Kafka-продюсер не отправляет каждое сообщение немедленно: он накапливает их в буфере и сбрасывает пачками. Если приложение остановить до сброса — накопленные сообщения пропадут.
Spring Boot при остановке закрывает автонастроенную фабрику продюсеров: DefaultKafkaProducerFactory.destroy() вызывает producer.close(Duration.ofSeconds(30)) — это значение physicalCloseTimeout по умолчанию. Буфер уходит на брокер, соединение закрывается. Если вы используете стандартный KafkaTemplate поверх этой фабрики — ничего дополнительно настраивать не нужно.
Если в проекте используется KafkaProducer напрямую (нестандартная ситуация), нужно явно задать destroyMethod:
@Configuration
public class RawProducerConfig {
private KafkaProducer<String, Object> producer;
@Bean(destroyMethod = "")
KafkaProducer<String, Object> rawProducer() {
this.producer = new KafkaProducer<>(producerProperties());
return this.producer;
}
@PreDestroy
public void closeProducer() {
producer.close(Duration.ofSeconds(15)); // явный flush с таймаутом
}
}
destroyMethod = "" отключает подбор метода закрытия: иначе Spring сам нашёл бы у бина close() без аргументов, а тот ждёт отправки буфера без ограничения по времени. Поэтому закрытие берут на себя — с таймаутом, и делает его @PreDestroy самого класса конфигурации, который держит ссылку на продюсер. Если оставить пустой destroyMethod и не добавить @PreDestroy, накопленные сообщения потеряются.
Транзакционный продюсер: остановка посреди транзакции
Если продюсер работает в транзакционном режиме (spring.kafka.producer.transaction-id-prefix задан), остановка посреди транзакции требует отдельного разговора — здесь ничего не теряется, но есть задержка, которая удивляет.
Что происходит при остановке. Незафиксированная транзакция остаётся незафиксированной. Сообщения физически записаны в журнал брокера, но помечены как часть незавершённой транзакции, и потребители с isolation.level=read_committed их не видят. Для них этих сообщений нет.
Что происходит дальше. Брокер ждёт завершения транзакции до истечения её таймаута (transaction.timeout.ms, по умолчанию 60 секунд у продюсера, ограничен настройкой брокера transaction.max.timeout.ms — 15 минут), после чего сам её прерывает. Пока этого не случилось, у потребителя с read_committed есть неприятный эффект: он не может двигаться дальше по партиции за точку начала незавершённой транзакции, потому что не знает, покажутся эти сообщения или нет. Внешне это выглядит как «потребление встало после выката», хотя ничего не сломано: просто ждём таймаута брошенной транзакции.
Как этого избежать. Три вещи.
Транзакции держат короткими: одна транзакция на одно сообщение или на маленькую пачку, а не на долгий обход. Долгая транзакция — это ровно тот случай, когда остановка оставляет её брошенной.
transaction.timeout.ms ставят соразмерно: значение по умолчанию (минута) для коротких транзакций избыточно велико, и при аварийной остановке ждать потребителям придётся именно минуту. Десять-пятнадцать секунд — разумный ориентир, если транзакции действительно короткие.
Остановку дают завершить: транзакционный продюсер закрывается вместе с контекстом, и если бюджета хватает, текущая транзакция успевает зафиксироваться или откатиться штатно — тогда никакого ожидания таймаута не возникает вовсе. То есть та же дисциплина бюджета, что и во всей фазе.
И отдельная деталь про связку «чтение и запись в одной транзакции» (когда сдвиг потребителя фиксируется той же транзакцией, что и отправка): при остановке посреди неё и сообщения, и сдвиг остаются незафиксированными. Это правильное поведение — после перезапуска пачка обработается заново целиком, — но оно требует идемпотентности на стороне получателя ваших сообщений, потому что первая попытка могла быть частично видна. Разбор — в статье про идемпотентность.
Как это проверить
Всё описанное проверяется одним прогоном на стенде, и без него настройки остаются предположением.
Подготовка. Тема с несколькими партициями, постоянный поток сообщений (простой продюсер в цикле или kafka-producer-perf-test), две-три реплики потребителя, включённая метрика отставания группы.
Что делаем. Под нагрузкой перезапускаем поды по одному:
kubectl rollout restart deployment/orders-consumer
kubectl rollout status deployment/orders-consumer
На что смотрим.
Отставание группы (kafka-consumer-groups --describe --group orders, колонка LAG, или метрика kafka_consumergroup_lag в системе метрик). Во время выката оно немного вырастет — это нормально; важно, что после выката оно возвращается к нулю, а не остаётся расти. Если растёт — потребление не восстановилось: обычно это затянувшееся перераспределение партиций.
Число повторно доставленных сообщений. Его видно по счётчику в самом обработчике: сколько раз он встретил уже обработанный идентификатор события. Ноль ожидать не стоит — пачка, не успевшая зафиксироваться, придёт заново, это штатное поведение; тревожит другое, когда повторов десятки процентов от потока: значит пачки слишком большие, а бюджет слишком мал.
Потерянные сообщения. Самое важное и самое простое: продюсер нумерует сообщения, потребитель пишет обработанные номера в базу, после прогона сверяем — пропусков быть не должно. Один пропуск означает, что где-то сдвиг зафиксирован раньше обработки; чаще всего это включённый вручную автокоммит.
Время остановки пода. kubectl get events или журнал: от начала остановки до исчезновения пода. Если оно равно бюджету, значит контейнер слушателя не успевает и его добивают — надо уменьшать max.poll.records.
Повторяют прогон после каждого изменения этих настроек. Полчаса на стенде против ночного разбора с потерянными сообщениями — обмен очевидный.
Глубже: rolling restart без ребаланса: group.instance.id и cooperative-stickyрасширенное
Всё выше про то, чтобы потребитель остановился аккуратно. Есть вторая цена остановки, которую платят соседи: когда под выходит из группы, брокер запускает ребаланс, партиции перераспределяются между оставшимися, и на время ребаланса обработка останавливается у всей группы. Rolling restart из пяти подов это пять ребалансов подряд плюс пять при возвращении, и десять пауз по несколько секунд на выкат.
Разница видна в третьей колонке: при обычном выходе встаёт вся группа, при постоянном идентификаторе соседи читают свои партиции дальше, а партиции ушедшего просто ждут его возвращения.
Статическое членство. Если у потребителя есть постоянный идентификатор, group.instance.id, брокер при его выходе не запускает ребаланс сразу, а ждёт session.timeout.ms: вернулся под с тем же идентификатором, получил те же партиции, и никто ничего не заметил. Идентификатор берут из имени пода, оно у StatefulSet постоянно, а у Deployment его заменяют порядковым номером реплики или именем пода с тем же префиксом:
spring:
kafka:
consumer:
properties:
group.instance.id: ${HOSTNAME}
session.timeout.ms: 60000
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor
session.timeout.ms поднимают до времени, за которое под успевает перезапуститься (минута), иначе окно ожидания закончится раньше и ребаланс всё равно случится; больше нескольких минут ставить нельзя, потому что столько же группа будет ждать под, который умер по-настоящему. Партиции остановившегося пода на это время просто не читаются, и отставание по ним растёт, что для большинства сервисов лучше остановки всей группы.
Кооперативный ребаланс. Стратегия по умолчанию отбирает у всех все партиции и раздаёт заново, и на это время стоит вся группа. CooperativeStickyAssignor перераспределяет только те партиции, которые действительно должны сменить владельца, а остальные продолжают читаться. Переход на неё с работающей группы делают в два выката: сначала список из двух стратегий, старой и кооперативной, потом одна кооперативная, иначе часть группы не поймёт протокол.
max.poll.interval.ms это третья настройка про то же: если обработчик не вернулся к poll за пять минут, брокер считает потребителя мёртвым и запускает ребаланс независимо от членства. Долгая пачка при остановке, о которой раздел про listener, упирается и в это: время обработки одной пачки обязано быть заметно меньше этого интервала, а не только меньше окна остановки.
С тремя настройками выкат перестаёт быть событием для группы: под уходит, его партиции ждут минуту, под возвращается и продолжает; соседи не останавливаются вовсе.
Коротко
- При остановке Spring просит
ConcurrentMessageListenerContainerостановиться — потребитель дорабатывает текущую пачку, коммитит offset и закрываетKafkaConsumer. spring.lifecycle.timeout-per-shutdown-phase: 20s— сколько Spring ждёт остановки контейнеров Kafka (по умолчанию 30 с); если listener не успевает, процесс завершится раньше, offset не зафиксируется и сообщения придут повторно.- Listener должен завершаться быстро — тяжёлые операции с повторами выносят в outbox, иначе приложение не уложится в timeout. Длительность обработки ограничивают два параметра:
max.poll.records(главный рычаг остановки, 500 по умолчанию — много) иmax.poll.interval.ms, за которым группа считает потребителя зависшим. ack-mode: BATCH— рекомендуемый режим: один коммит на пачку, при остановке поведение предсказуемо.enable.auto.commit: true— опасен: offset может зафиксироваться раньше реальной обработки, при shutdown сообщения теряются безвозвратно. Spring Kafka сам ставитenable.auto.commit=false: потеря сообщений возможна лишь если настройку включили руками, перенеся конфигурацию «голого» клиента.- Автонастроенная фабрика продюсеров сама сбрасывает буфер при остановке; для бина
KafkaProducerс пустымdestroyMethodнужен явныйclose(Duration)в@PreDestroy. - Выкат без ребаланса:
group.instance.idиз имени пода сsession.timeout.msпод время перезапуска,CooperativeStickyAssignorчерез два выката, обработка пачки заметно корочеmax.poll.interval.ms. spring.kafka.listener.immediate-stopпо умолчаниюfalse— контейнер дообрабатывает текущую пачку;trueберут только при дорогой обработке и гарантированной идемпотентности.- Остановка посреди транзакции продюсера ничего не теряет, но потребители с
read_committedне идут дальше до истеченияtransaction.timeout.ms— отсюда «потребление встало после выката». - Проверка на стенде: нагрузка плюс
rollout restart, смотреть возврат отставания к нулю, долю повторов, отсутствие пропусков в нумерации и время остановки пода.
Что почитать дальше
- JVM/Spring конфигурация — базовая настройка graceful shutdown в Spring Boot.
- Идемпотентность при остановке — как защититься от повторной обработки при replay.
- Бюджеты и наблюдаемость — Kafka в общем бюджете времени на shutdown.