1. Обзор
Приложения, основанные на MSA (микросервисной архитектуре), строятся путем разделения одной большой системы на несколько независимых сервисов. Каждый сервис выполняет разные функции и несет разные обязанности, однако для предоставления фактических сервисов необходимо продолжать передачу данных и совместное использование состояния между сервисами. Обмен сообщениями — один из методов взаимодействия между сервисами, а при необходимости асинхронной доставки данных можно использовать брокер сообщений.
Project A использует NATS JetStream — постоянную систему обмена сообщениями на базе NATS — для надежной и эффективной передачи данных датчиков и событий между несколькими сервисами на борту судна. NATS JetStream поддерживает не только асинхронную доставку сообщений между сервисами, но также их хранение, управление состоянием Consumer и повторную доставку, обеспечивая надежную обработку непрерывно генерируемых данных судна.
Однако во время фактической разработки и эксплуатации возникают различные ситуации, такие как принудительное выключение сервера судна и аварийное завершение приложения, и в таких условиях мы часто сталкивались с неожиданными сбоями, связанными с NATS JetStream. В этой статье на основе случаев сбоев, произошедших при эксплуатации NATS JetStream, мы анализируем основные причины и систематизируем процедуру диагностики для выявления и устранения подобных сбоев в будущем на основе значений состояния Consumer.
2. Возникновение сбоев обработки сообщений
После перезагрузки сервера вследствие аварийного выключения, например принудительного выключения, возникла ситуация, при которой Durable Consumer не мог корректно получать сообщения, хотя сообщения продолжали публиковаться в Stream.
При запуске это приложение проверяет наличие существующего Durable Consumer. Если Consumer уже существует, приложение повторно использует его вместо создания нового. Durable Consumer непрерывно хранит не только информацию о подписке, но и состояние потребления, включая Delivered Sequence, указывающий, до какого места были доставлены сообщения; Ack Floor, указывающий, до какого места подтверждения были завершены; состояние сообщений, которые еще не были подтверждены; и сообщения, предназначенные для повторной доставки. Поэтому при принудительном завершении работы сервера во время обработки сообщений подтверждение Ack для некоторых сообщений может не завершиться, либо состояние доставки и обработки Consumer может остаться незавершенным.
После перезагрузки сервера приложение продолжает повторно использовать существующий Durable Consumer и, следовательно, наследует состояние обработки сообщений, сохранявшееся непосредственно перед аварийным завершением. В обычных обстоятельствах необработанные сообщения должны доставляться повторно в соответствии с AckWait и политикой повторной доставки, после чего потребление сообщений должно продолжиться. Однако в зависимости от различных обстоятельств, например накопления сообщений Ack Pending или некорректного продолжения обработки состояния между Delivered Sequence и Ack Floor, доставка сообщений Consumer может замедлиться или остановиться.
3. Как анализировать значения состояния Consumer
При возникновении сбоя сначала проверьте информацию о состоянии Consumer и определите, на каком этапе произошел сбой.
Information for Consumer “stream” > “consumer”
Configuration:
Name: “consumer”
Pull Mode: true
Filter Subject: subject.>
Deliver Policy: All
Ack Policy: Explicit
Ack Wait: 30s
Maximum Deliveries: 3
Maximum Ack Pending: 1000
State:
Last Delivered Message:
Consumer sequence: 15230
Stream sequence: 185400
Acknowledgment Floor:
Consumer sequence: 14230
Stream sequence: 184400
OutStanding Acks: 1000
Redelivered Messages: 15
Unprocessed Messages: 3270
Waiting Pulls: 1
|
Показатель |
Значение |
Что проверить |
|---|---|---|
|
Последнее доставленное сообщение |
До какого места Consumer доставил сообщения |
Проверьте, не остановилась ли доставка на определенном Sequence на длительное время |
|
Граница подтверждений |
До какого места подтверждения выполнялись непрерывно |
Проверьте, не слишком ли велика разница с Last Delivered |
|
OutStanding Acks |
Количество доставленных, но еще не подтвержденных сообщений |
Проверьте, не достигло ли значение Maximum Ack Pending |
|
Необработанные сообщения |
Количество сообщений, которые еще не были доставлены Consumer |
Проверьте, продолжается ли потребление, если значение больше 0 |
|
Повторно доставленные сообщения |
Количество повторно доставленных сообщений |
Проверьте, не увеличивается ли оно аномально после перезагрузки сервера |
|
Ожидающие Pull-запросы |
Количество Pull-запросов, которые в данный момент запрашивают сообщения и ожидают их получения |
Если значение равно 0, проверьте, выполняется ли цикл Pull приложения |
-
Таблица 1. Основные значения состояния Consumer и их значение
1) Когда OutStanding Acks достигает MaxAckPending
OutStanding Acks — это количество сообщений, которые Consumer доставил приложению, но для которых еще не получил Ack. Когда это значение достигает настройки MaxAckPending, JetStream может ограничить доставку новых сообщений до тех пор, пока Acks не будут обработаны и не появится свободная емкость.
Maximum Ack Pending: 1000
OutStanding Acks: 1000
Unprocessed Messages: 3270
В приведённом выше примере в Stream всё ещё остаётся 3270 сообщений, но Consumer уже доставил сообщения до разрешённого предела Ack Pending. Поэтому, если приложение не обрабатывает Acks, дополнительные сообщения доставляться не будут, что может привести к ситуации, когда Publish извне выглядит нормально, но Consumer не получает сообщений.
2) Когда Pull-запросы не выполняются должным образом
Pull Consumer доставляет сообщения только тогда, когда приложение отправляет Pull-запросы с помощью fetch() или consume(). Поэтому, если сообщения остаются недоставленными, а Waiting Pulls остается равным 0, следует подозревать логику обработки Pull в приложении, а не сам Consumer.
Unprocessed Messages: 3270
OutStanding Acks: 0
Waiting Pulls: 0
В приведенном выше состоянии доставка сообщений не заблокирована из-за Acks, но отсутствуют Pull-запросы на получение сообщений. После перезагрузки сервера проверьте, корректно ли запустился Thread обработки Pull или цикл consume() и была ли должным образом восстановлена Subscription после повторного подключения к NATS.
3) Когда Last Delivered и Ack Floor не изменяются в течение длительного времени
Last Delivered указывает, до какого места Consumer доставил сообщения, а Ack Floor — до какого места подтверждения выполнялись непрерывно. Важно сосредоточиться не столько на разнице между этими двумя значениями, сколько на проверке того, не перестали ли оба значения изменяться с течением времени.
Last Delivered Stream Sequence: 185400
Ack Floor Stream Sequence: 185380
OutStanding Acks: 20
Unprocessed Messages: 3270
\n Если Last Delivered остается равным 185400 в течение нескольких секунд или минут, а Unprocessed Messages продолжает увеличиваться, поток потребления Consumer может быть остановлен. В этом случае одновременно проверьте Waiting Pulls, Redelivered Messages и журналы приложения, чтобы определить, связана ли проблема с Consumer или с обработкой в приложении.
В этих ситуациях сбоя потребление сообщений продолжалось нормально после удаления существующего Durable Consumer и перезапуска приложения для создания нового Consumer. Это подтвердило, что причиной была не Publisher и не сам Stream: на потребление сообщений повлияло состояние, сохраненное существующим Durable Consumer до аварийного завершения.
Впоследствии мы создали скрипт восстановления, который при перезапуске сервера удаляет и заново создает существующий Consumer, тем самым инициализируя состояние системы обмена сообщениями. Это помогло сократить количество сбоев, связанных с NATS JetStream на борту судна.
4. Заключение
В реальных условиях эксплуатации обработка сообщений в NATS JetStream может завершаться сбоем из-за непредвиденных ситуаций, таких как принудительное выключение сервера или аварийное завершение приложения. В частности, обработка сообщений является ключевой функцией, связанной с передачей данных между сервисами, поэтому при возникновении сбоя важно быстро определить его причину и принять соответствующие меры.
Поэтому при возникновении подобных сбоев в будущем необходимо быстро выявлять их причины и реагировать на них на основе основных значений состояния Consumer и методов анализа конкретных случаев, описанных в этой статье. Это позволит минимизировать время простоя сервисов и поддерживать стабильную эксплуатацию.
deeeneee