Это вторая статья серии «Погружение в NATS». Первая статья была про Core NATS — subjects, pub/sub, queue groups, request/reply — и заканчивалась одной важной оговоркой: Core ничего не хранит. Нет подписчика в момент публикации — сообщения больше нет. Это осознанная граница модели, а не недоработка, но она означает, что Core закрывает не все задачи, которые обычно решают брокеры сообщений.
JetStream — это ответ на вопрос «а как хранить». Он добавлен поверх того же сервера как отдельный слой: persistence, retention policies, consumers с ack и redelivery, дедупликация, KV и Object store. Именно JetStream в обзорной статье о выборе брокера закрывал задачи, которые обычно относят к Kafka и RabbitMQ. Здесь — подробно, с тем же принципом, что и в первой статье: если вы знаете Kafka или RabbitMQ, у вас есть вся нужная интуиция, нужно только сопоставить термины.
В статье
- Что такое JetStream и как включить
- Streams
- Retention policies
- Consumers
- Дедупликация и «exactly-once»
- KV и Object store
- Демонстрация
- Вывод
Что такое 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
--storage — file (диск, переживает рестарт) или memory (быстрее, но пропадает при падении процесса). file — выбор по умолчанию для всего, что должно хоть немного походить на надёжное хранилище.
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
Мост от 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
Важная оговорка, которая касается всех трёх: лимиты (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 минут», распределяя нагрузку на повторную обработку во времени.
дальнейшая судьба — по 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
Мост от 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-репликация — тема третьей статьи серии.
Комментарии