NATS из Go: nats.go, Core и JetStream паттерны

Практика на nats.go: подключение и reconnect, Drain и graceful shutdown, request/reply и queue groups, новый jetstream API с pull-consumer, ack и дедупликацией, KV — с оглядкой на идиомы клиентов Kafka и RabbitMQ

Пять первых статей серии «Погружение в NATS» разбирали модель: subjects и Core (ст. 1), streams и consumers в JetStream (ст. 2), кластеризацию, geo-топологию и производительность. Всё это — концепции, которые не зависят от языка клиента. Эта статья — первая из двух, которые переносят те же концепции в код: здесь — Go и nats.go, в следующей — Java и jnats. Паттерны намеренно называются одинаково в обеих статьях: подключение и устойчивость, Core pub/sub и queue groups, request/reply, JetStream, KV — чтобы при чтении седьмой статьи не пришлось заново разбираться, что к чему относится.

NATS из Go: nats.go, Core и JetStream паттерны

Все примеры ниже — рабочий код, а не фрагменты для иллюстрации: полный demo лежит в nats/clients/go в digital-cookbook и подключается к стендам 01-jetstream/02-cluster из предыдущих статей. Версия клиента — github.com/nats-io/nats.go v1.52.0, сервер — nats-server 2.12.x, как и во всей серии.

В статье

Подключение и устойчивость

Первое, что стоит знать про 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, поднять алерт.

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

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
Жизненный цикл соединения nats.go: connect → reconnect → drain → close

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 — в следующей статье серии.

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

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

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

Комментарии