JetStream: как NATS получает persistence, streams и гарантии доставки

Второй уровень NATS: streams и retention policies (limits/interest/workqueue), pull vs push consumers, ack и redelivery, дедупликация, KV и Object store — и чем это отличается от лога Kafka и очередей RabbitMQ

Это вторая статья серии «Погружение в NATS». Первая статья была про Core NATS — subjects, pub/sub, queue groups, request/reply — и заканчивалась одной важной оговоркой: Core ничего не хранит. Нет подписчика в момент публикации — сообщения больше нет. Это осознанная граница модели, а не недоработка, но она означает, что Core закрывает не все задачи, которые обычно решают брокеры сообщений.

JetStream: как NATS получает persistence, streams и гарантии доставки

JetStream — это ответ на вопрос «а как хранить». Он добавлен поверх того же сервера как отдельный слой: persistence, retention policies, consumers с ack и redelivery, дедупликация, KV и Object store. Именно JetStream в обзорной статье о выборе брокера закрывал задачи, которые обычно относят к Kafka и RabbitMQ. Здесь — подробно, с тем же принципом, что и в первой статье: если вы знаете Kafka или RabbitMQ, у вас есть вся нужная интуиция, нужно только сопоставить термины.

В статье

Что такое JetStream и как включить

JetStream — не отдельный процесс и не отдельный сервер, а подсистема того же nats-server, выключенная по умолчанию. Включается флагом -js (--jetstream) при старте:

nats-server -js -sd /data/jetstream

-sd (--store_dir) задаёт каталог на диске, куда JetStream пишет данные streams и consumers — метаданные и, если storage file, сами сообщения. Без -sd каталог хранения по умолчанию — подкаталог jetstream во временной директории ОС (/tmp/jetstream на Linux), что годится для локального эксперимента и не годится ни для чего больше: временный каталог может быть очищен ОС или потерян при перезапуске контейнера без volume.

В конфигурационном файле то же самое выглядит как отдельный блок:

jetstream {
  store_dir: /data/jetstream
  max_mem: 1G
  max_file: 100G
}

max_mem и max_file — верхние границы, которые сервер резервирует под JetStream в целом; конкретные streams дальше ограничиваются уже своими лимитами.

Важная деталь для 2.12.x: сервер теперь по умолчанию работает в «строгом режиме» — некорректные JetStream-запросы не просто логируются, а возвращаются клиенту как ошибка. Раньше часть таких ошибок можно было не заметить; теперь API строже проверяет входные данные, и это стоит учитывать при отладке скриптов и клиентского кода, написанного под более старые версии.

С точки зрения топологии на этом этапе ничего не меняется: один узел, один порт 4222 для клиентов, один порт для мониторинга. Кластеризация и RAFT-репликация streams между нодами — тема третьей статьи серии; здесь — один узел с локальным диском.

Streams

Stream — это именованный поток сообщений: JetStream подписывается на набор subjects (обычно с wildcard) и сохраняет всё, что в них публикуется, в указанном порядке. В отличие от Core NATS, где subject — просто адрес доставки, здесь subject становится ещё и критерием отбора сообщений в постоянное хранилище.

nats stream add ORDERS --subjects "orders.>" --storage file --retention limits --defaults

--storagefile (диск, переживает рестарт) или memory (быстрее, но пропадает при падении процесса). file — выбор по умолчанию для всего, что должно хоть немного походить на надёжное хранилище.

graph LR P["Publisher"] -->|"orders.eu.new"| S["Stream ORDERS
subjects: orders.>"] P2["Publisher"] -->|"orders.us.new"| S S --> C1["Consumer proc
filter: orders.eu.>"] S --> C2["Consumer audit
filter: orders.>"] style S fill:#f9f3e3,stroke:#8b7355 style C1 fill:#c9e4c5,stroke:#5b8a5e style C2 fill:#c9e4c5,stroke:#5b8a5e

graph LR
  P["Publisher"] -->|"orders.eu.new"| S["Stream ORDERS
subjects: orders.>"] P2["Publisher"] -->|"orders.us.new"| S S --> C1["Consumer proc
filter: orders.eu.>"] S --> C2["Consumer audit
filter: orders.>"] style S fill:#f9f3e3,stroke:#8b7355 style C1 fill:#c9e4c5,stroke:#5b8a5e style C2 fill:#c9e4c5,stroke:#5b8a5e
Stream ORDERS с фильтром orders.> и два consumer с разными подписками поверх одного stream

Мост от Kafka и RabbitMQ. Ближе всего по роли — topic-log в Kafka: и там, и там сообщения физически лежат в порядке публикации, и несколько независимых читателей могут пройти по этому логу с разной скоростью и разной позицией. Разница — в единице хранения. В Kafka единица — партиция топика, и параллелизм строится делением на партиции. В JetStream единица — stream, а внутренний параллелизм записи определяется не партициями, а количеством и конфигурацией consumers поверх одного и того же потока. От RabbitMQ stream отличается сильнее: в RabbitMQ единица хранения — очередь, обычно одна на потребителя или на логическую группу потребителей, и очередей заводят много под разные назначения. Один stream JetStream с фильтром orders.> по смыслу ближе к «одному большому логу заказов», поверх которого потом навешивают несколько consumers с разными подписками — это не то же самое, что заводить отдельную очередь на каждый сценарий обработки.

Retention policies

Retention policy — это то, что фактически определяет, ведёт себя ли stream как лог или как очередь. Их три:

  • limits (по умолчанию) — сообщения хранятся, пока не сработает один из лимитов: MaxMsgs, MaxBytes, MaxAge, MaxMsgsPerSubject. Ack потребителей на удаление не влияет вообще. Это прямой аналог retention в Kafka: лог живёт по времени/размеру, а не по факту прочтения.
  • interest — сообщение хранится, пока у него есть хотя бы один consumer, который его ещё не подтвердил; как только все существующие consumers его подтвердили, сообщение удаляется. Особенность: consumers должны существовать до публикации, иначе интереса к сообщению просто не будет.
  • workqueue — сообщение удаляется, как только его подтвердил любой consumer (при этом на один subject внутри stream может быть только один активный consumer — ограничение, которое JetStream проверяет при создании второго). Это прямой аналог классической очереди RabbitMQ: сообщение живёт, пока не обработано, и исчезает после ack.
graph TD subgraph L["limits — как лог Kafka"] L1["Сообщение хранится до MaxAge/MaxMsgs/MaxBytes"] L2["Ack потребителя не влияет на удаление"] end subgraph I["interest"] I1["Хранится, пока есть непрочитавший consumer"] I2["Все consumers подтвердили → удаление"] end subgraph W["workqueue — как очередь RabbitMQ"] W1["Один активный consumer на subject"] W2["Ack → немедленное удаление"] end style L fill:#f9f3e3,stroke:#8b7355 style I fill:#f9f3e3,stroke:#8b7355 style W fill:#c9e4c5,stroke:#5b8a5e

graph TD
  subgraph L["limits — как лог Kafka"]
    L1["Сообщение хранится до MaxAge/MaxMsgs/MaxBytes"]
    L2["Ack потребителя не влияет на удаление"]
  end
  subgraph I["interest"]
    I1["Хранится, пока есть непрочитавший consumer"]
    I2["Все consumers подтвердили → удаление"]
  end
  subgraph W["workqueue — как очередь RabbitMQ"]
    W1["Один активный consumer на subject"]
    W2["Ack → немедленное удаление"]
  end

  style L fill:#f9f3e3,stroke:#8b7355
  style I fill:#f9f3e3,stroke:#8b7355
  style W fill:#c9e4c5,stroke:#5b8a5e
Три retention policy: когда сообщение покидает stream

Важная оговорка, которая касается всех трёх: лимиты (MaxMsgs, MaxBytes, MaxAge) применяются поверх любой retention policy. Если на stream с workqueue задан MaxMsgs, старые неподтверждённые сообщения всё равно будут удалены при превышении лимита — retention policy определяет дополнительное условие удаления, а не отменяет базовые ограничения размера.

Практический вывод: limits — это «дайте мне лог с историей», workqueue — «дайте мне очередь задач, которая не должна расти», interest — промежуточный случай для нескольких параллельных читателей одного потока, ни один из которых не должен провоцировать неограниченный рост.

Consumers

Consumer — это именованный курсор поверх stream: он хранит позицию чтения, фильтр по subject, ack policy и параметры redelivery. Один stream может обслуживать множество consumers одновременно, каждый со своей независимой позицией.

Durable vs ephemeral. Durable consumer создаётся с явным именем и переживает рестарт клиента и сервера — позиция чтения сохраняется. Ephemeral не имеет постоянного имени, живёт в памяти сервера и удаляется, как только последняя подписка на него отваливается. Для сервисов, которым важно не потерять позицию между рестартами (а это почти всегда), нужен durable.

Pull vs push. Push consumer сам присылает сообщения в указанный subject — это ближе к push-модели RabbitMQ. Pull consumer, наоборот, ничего не присылает сам: клиент явно запрашивает батч сообщений вызовом вроде nats consumer next, и это сейчас рекомендуемый режим для новых проектов — он даёт клиенту полный контроль над темпом обработки и упрощает горизонтальное масштабирование воркеров.

nats consumer add ORDERS proc --pull --ack explicit --defaults
nats consumer next ORDERS proc

Мост от Kafka. Pull consumer в JetStream по духу ближе всего к poll() в Kafka consumer API: клиент явно приходит за очередной порцией сообщений, и именно клиент решает, когда и сколько запрашивать. Push consumer, наоборот, роднее с моделью RabbitMQ, где брокер сам толкает сообщения подписанному клиенту с учётом prefetch.

Ack policy — что сервер засчитывает подтверждением:

  • none — сообщение считается подтверждённым сразу при доставке, ack от клиента не нужен и не проверяется;
  • all — подтверждение одного сообщения автоматически подтверждает все более ранние в этом consumer;
  • explicit — подтверждать нужно каждое сообщение отдельно; это единственный режим, который включает точный контроль redelivery, и это выбор по умолчанию для всего, что похоже на обработку задач.

Redelivery. Если сообщение доставлено, но не подтверждено в течение AckWait, JetStream доставляет его снова — консьюмеру, который освободился, необязательно тому же самому. MaxDeliver ограничивает число попыток; после его исчерпания сообщение остаётся в stream (retention limits/interest) или требует ручного вмешательства (retention workqueue, где оно тоже не удаляется автоматически — специально, чтобы не терять данные молча). Backoff задаёт возрастающие интервалы между повторными попытками вместо фиксированного AckWait — то есть можно настроить что-то вроде «повторить через 5с, потом через 30с, потом через 5 минут», распределяя нагрузку на повторную обработку во времени.

sequenceDiagram participant C as Consumer (pull) participant S as JetStream C->>S: consumer next (запрос батча) S->>C: msg #1 (delivery attempt 1) Note over C: обработка не завершена, ack не отправлен Note over S: AckWait истёк S->>C: msg #1 (delivery attempt 2, redelivery) C->>S: ack Note over S: сообщение подтверждено,
дальнейшая судьба — по retention policy

sequenceDiagram
  participant C as Consumer (pull)
  participant S as JetStream

  C->>S: consumer next (запрос батча)
  S->>C: msg #1 (delivery attempt 1)
  Note over C: обработка не завершена, ack не отправлен
  Note over S: AckWait истёк
  S->>C: msg #1 (delivery attempt 2, redelivery)
  C->>S: ack
  Note over S: сообщение подтверждено,
дальнейшая судьба — по retention policy
Pull consumer: доставка, отсутствие ack, redelivery по AckWait

Мост от RabbitMQ. AckWait + MaxDeliver в JetStream закрывают ту же задачу, что ручной ack + prefetch + DLQ-политика в RabbitMQ, но иначе: в RabbitMQ переотправка происходит при разрыве соединения consumer или явном nack, а таймаута ожидания ack по умолчанию нет (это нужно строить отдельно, например через x-consumer-timeout на уровне политики). В JetStream таймаут встроен в саму модель consumer с самого начала — AckWait обязателен концептуально, даже если использовать значение по умолчанию.

Дедупликация и «exactly-once»

При публикации можно указать заголовок Nats-Msg-Id — клиентский уникальный идентификатор сообщения. Если JetStream получает публикацию с уже виденным Nats-Msg-Id в пределах окна дедупликации (Duplicate Window, по умолчанию 2 минуты), второе сообщение молча отбрасывается — в stream остаётся только первое.

nats pub orders.new '{"id":1}' -H Nats-Msg-Id:1
nats pub orders.new '{"id":1}' -H Nats-Msg-Id:1   # дубликат — не создаёт второй записи

Это решает конкретную и узкую проблему: publisher отправил сообщение, не получил подтверждение (например, оборвалась сеть), и на всякий случай повторяет отправку. Без дедупликации это дало бы дубль в stream; с Nats-Msg-Id — нет.

Важно не путать это с «exactly-once» в широком смысле. Дедупликация на публикации защищает от повторной записи в stream в пределах окна, но не гарантирует, что сообщение не будет обработано дважды где-то дальше — например, если consumer обработал сообщение, но не успел отправить ack до истечения AckWait, JetStream пришлёт его повторно, и это уже не дубликат с точки зрения дедупа (это то же самое сообщение, доставленное второй раз тому же или другому consumer). Поэтому практическая формула exactly-once в JetStream — это Nats-Msg-Id + окно дедупа на входе плюс идемпотентная обработка на выходе, а не встроенная гарантия «сервер сам разберётся». Ровно то же самое верно и для Kafka с его enable.idempotence/transactional producer, и для RabbitMQ, где надёжная обработка без дублей строится через тот же принцип идемпотентности потребителя — магии для точного однократного выполнения ни у одной из систем нет.

KV и Object store

Поверх streams JetStream предоставляет два готовых абстрактных API, о которых стоит знать, даже если детально они будут разобраны отдельно.

KV (key-value store) — это bucket с интерфейсом put/get/delete, материализованный как обычный stream с именем KV_<bucket> (по одному сообщению на ключ, с историей ревизий). Годится для конфигурации, feature flags, небольших состояний — везде, где раньше тянулся бы отдельный Redis или etcd ради простого key-value.

nats kv add config
nats kv put config feature.flag true
nats kv get config feature.flag

Object store — то же самое, но для файлов и блобов произвольного размера: под капотом объект режется на чанки и хранится как stream с чанками плюс метаданные. Полезно, когда нужно передать через NATS что-то крупнее одного сообщения — артефакт сборки, снапшот, файл конфигурации целиком — не поднимая отдельного S3-совместимого хранилища ради разовой задачи.

Оба API стоит рассматривать как удобную надстройку для частных случаев, а не замену stream/consumer модели — когда нужен полный контроль над retention, ack и redelivery, работают напрямую со stream.

Демонстрация

Рабочий стенд — nats/01-jetstream в digital-cookbook. Один узел NATS (образ nats:2.12) с включённым JetStream и volume под store dir.

git clone https://github.com/khorost-tech/digital-cookbook.git
cd digital-cookbook/nats/01-jetstream
docker compose up -d

Дальше в README демо — создание stream ORDERS с retention limits, pull-consumer с explicit ack, публикация с Nats-Msg-Id для проверки дедупликации, отдельный stream JOBS с retention workqueue, где сообщение физически исчезает из stream сразу после ack, и пара команд nats kv. Все команды и ожидаемый вывод — в README демо.

Вывод

JetStream — не замена Core NATS, а слой поверх него: тот же сервер, тот же протокол, но с явным включением через -js/jetstream {} и store dir на диске. Stream превращает subject-фильтр в постоянный поток сообщений; retention policy определяет, ведёт ли он себя как лог (limits), как очередь с разделяемым интересом (interest) или как классическая work-очередь, где сообщение исчезает после обработки (workqueue). Consumers — durable или ephemeral, pull или push — дают gap между «сообщение доставлено» и «сообщение обработано», закрытый через ack policy, AckWait и MaxDeliver. Дедупликация по Nats-Msg-Id защищает от повторной записи при retry публикации, но exactly-once в полном смысле всё равно требует идемпотентной обработки на стороне потребителя. KV и Object store — готовые надстройки для частных случаев, не более того.

Всё это по-прежнему работает на одном узле. Что произойдёт со streams и consumers, когда узлов становится три и включается RAFT-репликация — тема третьей статьи серии.

Документация

Обсуждение в Telegram

Присоединиться →

Комментарии