Пять первых статей серии «Погружение в NATS» разбирали модель: subjects и Core (ст. 1), streams и consumers в JetStream (ст. 2), кластеризацию, geo-топологию и производительность. Всё это — концепции, которые не зависят от языка клиента. Эта статья — первая из двух, которые переносят те же концепции в код: здесь — Go и nats.go, в следующей — Java и jnats. Паттерны намеренно называются одинаково в обеих статьях: подключение и устойчивость, Core pub/sub и queue groups, request/reply, JetStream, KV — чтобы при чтении седьмой статьи не пришлось заново разбираться, что к чему относится.
Все примеры ниже — рабочий код, а не фрагменты для иллюстрации: полный demo лежит в nats/clients/go в digital-cookbook и подключается к стендам 01-jetstream/02-cluster из предыдущих статей. Версия клиента — github.com/nats-io/nats.go v1.52.0, сервер — nats-server 2.12.x, как и во всей серии.
В статье
- Подключение и устойчивость
- Core: pub/sub и queue groups
- Request/Reply
- JetStream: новый пакет
jetstream - KV
- Продакшн-советы
- Вывод
Подключение и устойчивость
Первое, что стоит знать про nats.go: клиент сам умеет переподключаться, и по умолчанию делает это разумно — но пара опций всё равно нужна почти всегда, а без коллбэков процесс молчит о переходах состояния, что в проде неприемлемо.
opts := []nats.Option{
nats.Name(name),
// RetryOnFailedConnect — не завершать nats.Connect ошибкой, если сервер
// временно недоступен в момент старта: клиент уйдёт в фоновый reconnect-
// цикл и подключится, как только сервер станет доступен.
nats.RetryOnFailedConnect(true),
// MaxReconnects(-1) — не ограничивать число попыток переподключения.
// По умолчанию клиент после исчерпания лимита попыток считает соединение
// окончательно потерянным и вызывает ClosedHandler — для долгоживущего
// сервиса это обычно не то поведение, которое нужно.
nats.MaxReconnects(-1),
nats.ReconnectWait(1 * time.Second),
nats.DisconnectErrHandler(func(_ *nats.Conn, err error) {
log.Printf("[%s] disconnected: %v", name, err)
}),
nats.ReconnectHandler(func(nc *nats.Conn) {
log.Printf("[%s] reconnected to %s", name, nc.ConnectedUrl())
}),
nats.ClosedHandler(func(_ *nats.Conn) {
log.Printf("[%s] connection closed", name)
}),
}
nc, err := nats.Connect(url, opts...)
Три коллбэка покрывают три разных события, и путать их не стоит: DisconnectErrHandler вызывается при потере соединения (клиент уже ушёл в фоновый reconnect-цикл), ReconnectHandler — когда соединение восстановлено, ClosedHandler — когда клиент прекратил попытки насовсем (обычно после явного Close() или, при ограниченном MaxReconnects, после исчерпания лимита). Если ClosedHandler сработал не по вашей команде — это сигнал, что сервис остался без брокера и это нужно как-то обработать: перезапуститься, уйти в circuit breaker, поднять алерт.
DisconnectErrHandler Reconnecting --> Connected: ReconnectHandler Reconnecting --> Reconnecting: ReconnectWait между попытками Connected --> Draining: Drain() Draining --> Closed: подписки и исходящие
дожиты, соединение закрыто Connected --> Closed: Close() Reconnecting --> Closed: MaxReconnects исчерпан Closed --> [*]: ClosedHandler
stateDiagram-v2
[*] --> Connecting: nats.Connect()
Connecting --> Connected: успех
Connecting --> Connecting: RetryOnFailedConnect(true)
Connected --> Reconnecting: обрыв связи
DisconnectErrHandler
Reconnecting --> Connected: ReconnectHandler
Reconnecting --> Reconnecting: ReconnectWait между попытками
Connected --> Draining: Drain()
Draining --> Closed: подписки и исходящие
дожиты, соединение закрыто
Connected --> Closed: Close()
Reconnecting --> Closed: MaxReconnects исчерпан
Closed --> [*]: ClosedHandler
Drain() против Close(). Close() рвёт соединение немедленно: недошедшие исходящие публикации и незавершённые обработчики подписок теряются. Drain() — это graceful shutdown в терминах pub/sub: соединение перестаёт принимать новые сообщения на подписки, ждёт завершения уже запущенных обработчиков, досылает буферизованные исходящие публикации — и только после этого закрывается само. Тот же принцип, что graceful shutdown у HTTP-сервера (дождаться in-flight запросов), просто перенесённый на подписки и публикации.
func Shutdown(ctx context.Context, nc *nats.Conn) {
done := make(chan struct{})
go func() {
if err := nc.Drain(); err != nil {
log.Printf("drain: %v", err)
}
close(done)
}()
select {
case <-done:
case <-ctx.Done():
log.Printf("drain: не завершился до дедлайна, закрываем принудительно")
nc.Close()
}
}
Этот хелпер (вместе с Connect) вынесен в internal/natsconn и переиспользуется всеми примерами demo — ниже в статье он не повторяется, только вызывается.
Мост от RabbitMQ и Kafka. В amqp091-go эквивалента Drain() нет — там принято закрывать Channel, затем Connection, и приложение само отвечает за то, чтобы дождаться in-flight обработчиков до этого момента (обычно через sync.WaitGroup вокруг горутин-консьюмеров). В franz-go/sarama для Kafka похожая идея есть в Close() консьюмер-группы — он тоже ждёт завершения текущего цикла poll/обработки перед выходом, но именование и детали API другие. Идея «дождаться, не оборвать» универсальна для всех трёх экосистем; nats.go просто даёт ей отдельное имя и метод.
Core: pub/sub и queue groups
nats.go даёт два способа читать сообщения: async subscribe — сервер сам вызывает переданный обработчик в отдельной горутине на каждое сообщение, и sync subscribe — клиент явно забирает следующее сообщение вызовом NextMsg.
// Async — обычный выбор для долгоживущих подписчиков.
subAsync, err := nc.Subscribe("orders.>", func(msg *nats.Msg) {
log.Printf("[async] %s: %s", msg.Subject, string(msg.Data))
})
// Sync — клиент сам решает, когда забрать следующее сообщение.
subSync, err := nc.SubscribeSync("orders.sync.>")
// ...
msg, err := subSync.NextMsg(2 * time.Second)
Sync-вариант пригождается там, где обработка должна идти строго последовательно в известной точке кода — например, в тестовом сценарии или CLI-инструменте, где не нужен отдельный обработчик в фоне и важнее явный контроль над порядком вызовов, чем параллелизм.
Queue groups — тот же встроенный механизм competing consumers, что и в nats CLI (ст. 1): несколько подписчиков с одинаковым именем группы делят входящий поток, каждое сообщение достаётся ровно одному участнику.
for i := 1; i <= 2; i++ {
workerID := i
sub, err := nc.QueueSubscribe("jobs.*", "workers", func(msg *nats.Msg) {
log.Printf("[queue worker-%d] %s: %s", workerID, msg.Subject, string(msg.Data))
})
// ...
}
Прогон трёх публикаций на jobs.a против двух воркеров в одной группе workers в demo cmd/core-pubsub подтверждает то же поведение, что было видно через nats CLI: сервер распределяет сообщения, а не рассылает их всем.
Одна деталь, о которой легко забыть при переходе с синхронных клиентов: Publish буферизует запись на стороне клиента, а не отправляет её на сервер немедленно. nc.Flush() дожидается, что буфер реально ушёл — полезно перед тем, как завершать короткоживущий процесс (CLI-утилиту, скрипт, demo), где иначе можно выйти раньше, чем данные покинут клиент.
Мост от RabbitMQ и Kafka. Subscribe с колбэком ближе всего к тому, как amqp091-go регистрирует consumer через Channel.Consume — оба отдают обработку сообщений в фоновый канал/горутину. SubscribeSync/NextMsg больше похож на явный poll() в franz-go/sarama — клиент сам решает, когда идти за следующим сообщением, только в Core NATS это работает поверх подписки, а не consumer group с офсетами. QueueSubscribe закрывает ту же задачу, что несколько consumer на одной очереди RabbitMQ или consumer group в Kafka, но в API nats.go это один явный параметр вызова, а не побочный эффект топологии очередей или конфигурации группы.
Request/Reply
Request/reply в nats.go — не паттерн поверх pub/sub, а отдельный вызов, использующий встроенный в протокол NATS inbox-механизм (ст. 1). Ответчик — обычная подписка, обычно в queue group, чтобы несколько экземпляров сервиса могли делить нагрузку:
const subject = "service.time"
sub, err := nc.QueueSubscribe(subject, "time-service", func(msg *nats.Msg) {
now := time.Now().Format(time.RFC3339)
if err := msg.Respond([]byte(now)); err != nil {
log.Printf("responder: ответ не отправлен: %v", err)
}
})
Клиент-requester задаёт таймаут явно — через time.Duration или через context, если вызов уже встроен в код с контекстом (HTTP-хендлер, gRPC-метод):
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
msg, err := nc.RequestWithContext(ctx, subject, nil)
switch {
case errors.Is(err, nats.ErrTimeout), errors.Is(err, context.DeadlineExceeded):
return errors.New("нет ответа — responder недоступен или перегружен")
case err != nil:
return err
}
nc.Request(subj, data, timeout) с time.Duration делает то же самое без контекста — выбор между двумя вариантами обычно диктует то, пронизан ли остальной код context.Context или нет; функционально они эквивалентны, разница в том, что RequestWithContext прокидывает отмену вызывающего контекста, а не только фиксированный таймаут.
Мост от RabbitMQ и Kafka. В amqp091-go request/reply собирается руками: временная эксклюзивная очередь (или direct reply-to), поле CorrelationId, сопоставление ответа на стороне клиента. В Kafka встроенного request/reply нет вовсе — обычно два независимых топика (запросов и ответов) плюс сопоставление по ключу, с задержкой заметно выше, чем у синхронного RPC. В nats.go это два вызова: QueueSubscribe + msg.Respond на одной стороне, RequestWithContext — на другой; inbox и сопоставление ответа целиком на стороне клиентской библиотеки и сервера.
JetStream: новый пакет jetstream
Важная развилка API, которую стоит проговорить явно: у nats.go есть два способа работать с JetStream. Старый — nc.JetStream(), возвращающий nats.JetStreamContext из пакета nats — всё ещё работает, но не развивается. Новый — отдельный пакет github.com/nats-io/nats.go/jetstream, с context.Context в каждом сетевом вызове и явным разделением управления (streams, consumers) и публикации/потребления на отдельные интерфейсы. Все примеры ниже — только новый пакет; он используется в тексте статьи умышленно однозначно, чтобы не смешивать два API.
import "github.com/nats-io/nats.go/jetstream"
js, err := jetstream.New(nc)
Stream и durable pull-consumer с explicit ack создаются идемпотентно — CreateOrUpdateStream/CreateOrUpdateConsumer не упадут, если ресурс уже существует с такой же конфигурацией:
stream, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: "ORDERS",
Subjects: []string{"orders.>"},
Storage: jetstream.FileStorage,
})
consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
Durable: "proc",
AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 10 * time.Second,
MaxDeliver: 5,
FilterSubject: "orders.>",
})
Дедупликация при публикации. Заголовок Nats-Msg-Id (ст. 2) передаётся опцией jetstream.WithMsgID, а не полем структуры — так проще добавлять другие опции публикации (например, jetstream.WithExpectStream) без раздувания сигнатуры:
ack1, err := js.Publish(ctx, "orders.new", payload, jetstream.WithMsgID("order-1"))
// ack1.Duplicate == false
ack2, err := js.Publish(ctx, "orders.new", payload, jetstream.WithMsgID("order-1"))
// ack2.Duplicate == true — сервер отбросил повтор, seq тот же
Прогон cmd/jetstream-pull против живого nats/01-jetstream печатает это напрямую:
publish #1: seq=1, duplicate=false
publish #2 (тот же MsgId): seq=1, duplicate=true
Кроме синхронного Publish, есть асинхронный PublishAsync — он возвращает PubAckFuture сразу, не дожидаясь подтверждения сервера, а результат разбирается позже через каналы Ok()/Err(). Полезно для высокого throughput публикации, когда ждать ack на каждое сообщение по очереди слишком дорого:
ackF, err := js.PublishAsync("orders.new", payload)
select {
case ack := <-ackF.Ok():
log.Printf("published: seq %d", ack.Sequence)
case err := <-ackF.Err():
log.Printf("publish async failed: %v", err)
}
(Если в коде или примерах где-то встречается имя AckAsync — это ошибка: правильное имя метода — PublishAsync, ack асинхронно приходит именно на публикацию, а не отдельным вызовом.)
Pull-consumer: Fetch и Consume. У нового пакета — два способа читать. Fetch забирает фиксированный батч и возвращает канал сообщений, которые нужно разобрать самостоятельно — подходит для одноразового чтения или мелких пакетных задач:
batch, err := consumer.Fetch(1, jetstream.FetchMaxWait(5*time.Second))
for msg := range batch.Messages() {
meta, _ := msg.Metadata()
log.Printf("получено: %s (попытка доставки: %d)", msg.Subject(), meta.NumDelivered)
if err := msg.Ack(); err != nil {
log.Printf("ack: %v", err)
}
}
Consume — для долгоживущего воркера: регистрирует колбэк, который сервер вызывает на каждое новое сообщение, а управление pull-запросами под капотом (клиент сам приходит за новым батчем, когда предыдущий обработан) библиотека берёт на себя:
cc, err := consumer.Consume(func(msg jetstream.Msg) {
// обработка
msg.Ack()
})
defer cc.Stop()
Демо cmd/jetstream-pull использует Fetch, чтобы процесс завершался предсказуемо после одного сообщения — в реальном сервисе, который живёт постоянно, Consume обычно удобнее.
Ack — не единственный вариант ответа. msg.Ack() подтверждает обработку. msg.Nak()/msg.NakWithDelay(d) — явный отказ: сообщение будет доставлено снова, опционально не раньше, чем через d. msg.InProgress() продлевает AckWait, не завершая обработку, — полезно для длинных задач, где AckWait короче фактического времени обработки. msg.Term() — «не пытаться больше», сообщение помечается как обработанное без успеха и redelivery прекращается, даже если MaxDeliver не исчерпан; годится для сообщений, которые в принципе невозможно обработать (битый payload, например) — не тратить MaxDeliver попыток впустую.
Мост от RabbitMQ и Kafka. CreateOrUpdateStream/CreateOrUpdateConsumer в JetStream закрывают ту же задачу, что декларация exchange/queue в amqp091-go — идемпотентное объявление топологии при старте сервиса. Fetch по духу ближе всего к poll() в franz-go/sarama: клиент явно приходит за очередной порцией, и именно он решает темп; Consume, наоборот, роднее с push-моделью RabbitMQ, где брокер сам присылает сообщения подписанному обработчику. Nak/Term не имеют точного аналога в терминах RabbitMQ (там это nack с requeue=true/false) — по смыслу Nak близок к nack(requeue=true), а Term — к nack(requeue=false): сообщение перестаёт переотправляться этому consumer’у и при этом не остаётся «неподтверждённым». Но отдельной DLQ, как очереди в RabbitMQ, в JetStream из коробки нет: Term только останавливает доставку. «Мёртвую» маршрутизацию собирают поверх advisory — подписываются на $JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES (или terminated-advisory) и по номеру достают исходное сообщение из stream, чтобы переложить или разобрать вручную. То есть DLQ здесь — прикладной механизм над advisory и чтением из stream, а не встроенная очередь.
KV
KV — API поверх JetStream (ст. 2), но со своим набором методов, а не через Publish/Fetch напрямую:
kv, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
Bucket: "config",
})
rev, err := kv.Put(ctx, "feature.flag", []byte("true"))
entry, err := kv.Get(ctx, "feature.flag")
// entry.Value(), entry.Revision()
watcher, err := kv.Watch(ctx, "feature.flag")
defer watcher.Stop()
for entry := range watcher.Updates() {
if entry == nil {
continue // разделитель между историей и живыми событиями
}
log.Printf("%s = %s (revision %d)", entry.Key(), string(entry.Value()), entry.Revision())
}
Практическая деталь, которая иначе стоит времени на отладку: Watch сразу присылает текущее значение ключа первым событием (это ожидаемо, не баг), а границу между «довычиткой истории» и «живыми» обновлениями отмечает nil-записью в канале — код выше её просто пропускает.
Продакшн-советы
Контексты и таймауты — везде, где есть сеть. Все вызовы нового jetstream-пакета принимают context.Context первым аргументом — используйте его: context.WithTimeout на публикацию и на операции управления streams/consumers предотвращает зависание сервиса, если JetStream временно недоступен или перегружен. Для Core RequestWithContext — то же самое для request/reply.
Redelivery — не исключение, а нормальный путь. AckWait и MaxDeliver должны быть осознанным выбором, а не значениями по умолчанию, скопированными из документации. Слишком короткий AckWait относительно реального времени обработки — гарантированные ложные redelivery и дублирующая работа; слишком длинный — задержка перед тем, как зависшее сообщение попадёт к другому воркеру. msg.InProgress() — способ сказать «я жив, просто это долго» без увеличения AckWait глобально для consumer.
Идемпотентность обработчика — обязательна, а не опциональна. Nats-Msg-Id защищает от дублирующей записи в stream при повторной публикации, но не от повторной доставки уже записанного сообщения — если consumer не успел подтвердить сообщение до истечения AckWait, JetStream пришлёт то же сообщение снова, и с точки зрения дедупа это не дубликат. Единственный надёжный способ закрыть это — идемпотентный обработчик на стороне consumer (по бизнес-ключу сообщения, не по факту доставки). Это верно для JetStream, Kafka с enable.idempotence и RabbitMQ в равной степени — ни одна из систем не даёт exactly-once обработки бесплатно.
Drain() в graceful shutdown сервиса, не только в demo. Подписывайтесь на SIGTERM/SIGINT и вызывайте Drain() с разумным дедлайном (nats.DrainTimeout(), по умолчанию 30 секунд) — это тот же принцип, что graceful shutdown HTTP-сервера, и его стоит применять одинаково последовательно.
Вывод
nats.go даёт готовые ответы на вопросы, которые в других экосистемах messaging приходится решать руками: reconnect с коллбэками из коробки, Drain() как отдельный примитив graceful shutdown, request/reply без временных очередей, dedup на публикации через один заголовок. Новый пакет jetstream — не косметическое обновление, а более строгий API: context.Context в каждом сетевом вызове, явное разделение управления ресурсами и потребления сообщений, Fetch/Consume как два явных режима работы pull-consumer вместо одного смешанного. Ничего из этого не отменяет базовых уроков ст. 1–2: at-most-once в Core, идемпотентность обработчика в JetStream — код лишь делает эти правила исполнимыми.
Тот же набор паттернов — на Java и jnats — в следующей статье серии.
Комментарии