NATS из Java: jnats, Core и JetStream паттерны

Те же паттерны на jnats: Connection и Dispatcher, request/reply, JetStream с pull-consumer через ConsumerContext, ack и redelivery, KV и graceful drain — зеркально к Go-статье серии

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

NATS из Java: jnats, Core и JetStream паттерны

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

В статье

Подключение

jnats собирает конфигурацию через Options.Builder — по духу тот же паттерн, что []nats.Option в nats.go, только это один объект-строитель, а не список функциональных опций:

Options options = new Options.Builder()
        .server(DEFAULT_URL)
        .connectionName(name)

        // maxReconnects(-1) — не ограничивать число попыток
        // переподключения. По умолчанию клиент после исчерпания
        // лимита попыток считает соединение окончательно потерянным
        // и переходит в CLOSED — для долгоживущего сервиса это
        // обычно не то поведение, которое нужно.
        .maxReconnects(-1)
        .reconnectWait(Duration.ofSeconds(1))

        .connectionListener((conn, type) ->
                System.out.printf("[%s] connection event: %s%n", name, type))
        .errorListener(new ErrorListener() {
            @Override
            public void errorOccurred(Connection conn, String error) {
                System.out.printf("[%s] error: %s%n", name, error);
            }

            @Override
            public void exceptionOccurred(Connection conn, Exception exp) {
                System.out.printf("[%s] exception: %s%n", name, exp);
            }
        })
        .build();

Connection nc = Nats.connect(options);

Разница с Go в самой форме коллбэков заметна сразу. В nats.go три отдельных хендлера — DisconnectErrHandler, ReconnectHandler, ClosedHandler — по одному на переход состояния. В jnats один ConnectionListener с единственным методом connectionEvent(Connection, Events), и все переходы приходят через него, различаясь значением Events: CONNECTED, DISCONNECTED, RECONNECTED, CLOSED, а также RESUBSCRIBED, DISCOVERED_SERVERS, LAME_DUCK, которым в Go-клиенте нет прямого аналога в виде отдельного callback (частично покрываются DiscoveredServersHandler, но не единым списком событий). ErrorListener — второй, отдельный интерфейс, покрывающий то, что в nats.go уходит в AsyncErrorHandler: ошибки протокола, slow consumer, discard сообщений.

stateDiagram-v2 [*] --> Connecting: Nats.connect() Connecting --> Connected: успех Connecting --> Connecting: reconnect при недоступности сервера Connected --> Reconnecting: обрыв связи
ConnectionListener: DISCONNECTED Reconnecting --> Connected: ConnectionListener: RECONNECTED Reconnecting --> Reconnecting: reconnectWait между попытками Connected --> Draining: drain() Draining --> Closed: подписки и исходящие
дожиты, соединение закрыто Connected --> Closed: close() Reconnecting --> Closed: maxReconnects исчерпан Closed --> [*]: ConnectionListener: CLOSED

stateDiagram-v2
    [*] --> Connecting: Nats.connect()
    Connecting --> Connected: успех
    Connecting --> Connecting: reconnect при недоступности сервера

    Connected --> Reconnecting: обрыв связи
ConnectionListener: DISCONNECTED Reconnecting --> Connected: ConnectionListener: RECONNECTED Reconnecting --> Reconnecting: reconnectWait между попытками Connected --> Draining: drain() Draining --> Closed: подписки и исходящие
дожиты, соединение закрыто Connected --> Closed: close() Reconnecting --> Closed: maxReconnects исчерпан Closed --> [*]: ConnectionListener: CLOSED
Жизненный цикл соединения (общий для nats.go и jnats): connect → reconnect → drain → close

drain() против close() — та же пара понятий, что в Go, только drain() в jnats возвращает CompletableFuture<Boolean>, а не блокирует вызывающего напрямую:

try (Connection nc = Nats.connect(options)) {
    // ... публикации и подписки ...

    nc.drain(Duration.ofSeconds(5)).get(10, TimeUnit.SECONDS);
}

Смысл идентичен Go-версии: close() (здесь — неявно через try-with-resources, если drain() не вызван явно до выхода из блока) рвёт соединение немедленно, недошедшие исходящие и незавершённые обработчики подписок теряются. drain() переводит соединение в режим «дожить»: перестаёт принимать новые сообщения на подписки, ждёт завершения уже запущенных обработчиков, досылает буферизованные исходящие публикации — и только после этого закрывается само. CompletableFuture.get(timeout, unit) здесь играет ту же роль, что context.Context с дедлайном вокруг Drain() в Go-хелпере Shutdown(): если drain() не успел за отведённое время, get бросит TimeoutException, и это тот сигнал, на котором стоит принудительно звать close().

Общая конфигурация подключения вынесена в NatsConn и переиспользуется всеми примерами demo — ниже в статье она не повторяется, только вызывается через NatsConn.options(name).

Core: pub/sub и queue groups

В jnats, как и в nats.go, два способа читать сообщения: async через фоновый обработчик и sync через явный вызов. Разница в том, как оформлен async-путь: в jnats для этого есть отдельная сущность — Dispatcher.

// Async — Dispatcher создаёт обработчик, который сервер вызывает в
// потоке диспетчера на каждое входящее сообщение. Один Dispatcher может
// обслуживать несколько подписок на общем потоке.
Dispatcher dispatcher = nc.createDispatcher(msg ->
        System.out.printf("[async] %s: %s%n", msg.getSubject(),
                new String(msg.getData(), StandardCharsets.UTF_8)));
dispatcher.subscribe("orders.>");

// Sync — клиент сам забирает сообщение вызовом nextMessage(timeout),
// без фонового обработчика.
Subscription subSync = nc.subscribe("orders.sync.>");
Message msg = subSync.nextMessage(Duration.ofSeconds(2));

Dispatcher — не просто способ подписаться, а единица управления потоком выполнения: Connection.createDispatcher(handler) создаёт один диспетчер с одним потоком, и на этот же диспетчер можно повесить несколько subscribe() — все они будут обрабатываться последовательно в этом потоке. В nats.go, для сравнения, у каждой подписки своя горутина по умолчанию, и явной концепции «общий поток на несколько подписок» нет — если нужно ограничить параллелизм, это делается пулом или каналом поверх обработчиков. Sync-путь (subscribe без обработчика + nextMessage) устроен одинаково в обоих клиентах и пригождается там же: обработка должна идти строго последовательно в известной точке кода, а не в фоне.

Queue groups — тот же встроенный механизм competing consumers, что и в nats CLI (ст. 1), только вызов идёт через Dispatcher:

String queue = "workers";
for (int i = 1; i <= 2; i++) {
    int workerId = i;
    dispatcher.subscribe("jobs.*", queue, msg ->
            System.out.printf("[queue worker-%d] %s: %s%n", workerId, msg.getSubject(),
                    new String(msg.getData(), StandardCharsets.UTF_8)));
}

Прогон трёх публикаций на jobs.a против двух воркеров в одной группе workers в demo CorePubSub подтверждает то же поведение, что и в Go-демо: сервер распределяет сообщения между участниками группы, а не рассылает их всем. Connection.subscribe(subject, queue) даёт то же самое для sync-подписки, если async через Dispatcher не нужен.

Как и в nats.go, publish() буферизует запись на стороне клиента — Connection.flush(timeout) дожидается, что буфер реально ушёл на сервер, что полезно перед завершением короткоживущего процесса.

Мост от Kafka и RabbitMQ. Dispatcher с обработчиком ближе всего к тому, как Kafka Java client регистрирует ConsumerRebalanceListener поверх poll()-цикла — только в jnats цикл чтения скрыт внутри клиента, а не в коде приложения; идея разделяемого потока на несколько подписок больше напоминает SimpleMessageListenerContainer в Spring AMQP, где один контейнер тоже может слушать несколько очередей. Sync nextMessage() ближе к явному poll(Duration) Kafka-консьюмера — тем же способом клиент сам управляет темпом чтения. Dispatcher.subscribe(subject, queue) закрывает ту же задачу, что consumer group в Kafka или несколько consumer на одной очереди RabbitMQ, но здесь это один параметр вызова, а не побочный эффект конфигурации топологии.

Request/Reply

Request/reply в jnats — как и в Go — использует встроенный в протокол NATS inbox-механизм (ст. 1), а не собирается поверх pub/sub вручную. Ответчик — обычная подписка через Dispatcher, в queue group, чтобы несколько экземпляров сервиса делили нагрузку:

private static final String SUBJECT = "service.time";

Dispatcher dispatcher = nc.createDispatcher(msg -> {
    String now = OffsetDateTime.now().toString();
    nc.publish(msg.getReplyTo(), now.getBytes(StandardCharsets.UTF_8));
});
dispatcher.subscribe(SUBJECT, "time-service");

Обратите внимание: ответ отправляется обычным Connection.publish(replyTo, body) на msg.getReplyTo(), а не отдельным методом вроде msg.Respond() в nats.go — в jnats эта пара методов не выделена в отдельный API, getReplyTo() — это просто subject inbox’а, куда отвечает publish.

Клиент-requester получает CompletableFuture<Message> и сам решает, как задавать таймаут — через get(timeout, unit):

CompletableFuture<Message> future = nc.request(SUBJECT, null);
try {
    Message msg = future.get(2, TimeUnit.SECONDS);
    System.out.printf("requester: получен ответ: %s%n",
            new String(msg.getData(), StandardCharsets.UTF_8));
} catch (TimeoutException e) {
    System.out.println("requester: нет ответа за 2s — responder недоступен или перегружен");
} catch (ExecutionException e) {
    System.out.printf("requester: request завершился ошибкой: %s%n", e.getCause());
}

Это заметное отличие от Go: там таймаут — параметр самого вызова (nc.Request(subj, data, timeout)) или контекста (RequestWithContext), и оба способа эквивалентны по возможностям. В jnats Connection.request(subject, body) всегда возвращает CompletableFuture без встроенного таймаута — таймаут навязывается снаружи, через get(timeout, unit), что даёт больше гибкости (можно, например, скомпоновать future с другими асинхронными вызовами через thenCompose/thenApply), но требует не забыть его вызвать: без таймаута ожидание ответа, которого не будет, зависнет навсегда. Есть и блокирующий вариант Connection.request(subject, body, timeout), который возвращает Message напрямую и бросает InterruptedException — по духу он ближе к nc.Request из Go, но не даёт отмены по контексту, только фиксированный таймаут.

Мост от Kafka и RabbitMQ. Ни у Kafka Java client, ни у amqp-client/Spring AMQP нет встроенного request/reply — в Kafka это обычно два независимых топика (запросов и ответов) с сопоставлением по ключу, в RabbitMQ — временная эксклюзивная очередь или directReplyTo плюс correlationId, который приложение сопоставляет само. В jnats это два метода: Dispatcher.subscribe + publish(getReplyTo(), ...) на одной стороне, request() — на другой; inbox и сопоставление ответа целиком на стороне клиентской библиотеки и сервера, задержка сопоставима с прямым RPC.

JetStream

jnats даёт доступ к JetStream через два интерфейса: JetStream (публикация и подписка на сообщения) и JetStreamManagement (управление streams и consumers — их создание, изменение, удаление). Это разделение прямо соответствует разделению в новом пакете jetstream из Go-статьи, только оформлено как два отдельных интерфейса вместо одного объекта с методами управления и потребления вперемешку:

JetStreamManagement jsm = nc.jetStreamManagement();
JetStream js = nc.jetStream();

Stream создаётся через JetStreamManagement.addStream. Важная деталь, которой в jnats в явном виде нет: в отличие от Go, где CreateOrUpdateStream/CreateOrUpdateConsumer идемпотентны из коробки, у JetStreamManagement нет единого «создать-или-обновить» для streams — только addStream/updateStream по отдельности. Consumers в этом смысле удобнее: для них есть addOrUpdateConsumer. Идемпотентность для stream в demo реализована вручную — через перехват кода ошибки «уже существует»:

private static StreamInfo createOrGetStream(JetStreamManagement jsm) throws Exception {
    try {
        return jsm.addStream(StreamConfiguration.builder()
                .name(STREAM_NAME)
                .subjects(STREAM_SUBJECT)
                .storageType(StorageType.File)
                .build());
    } catch (JetStreamApiException e) {
        if (e.getApiErrorCode() == 10058) { // stream name already in use
            return jsm.getStreamInfo(STREAM_NAME);
        }
        throw e;
    }
}
jsm.addOrUpdateConsumer(STREAM_NAME, ConsumerConfiguration.builder()
        .durable(CONSUMER_NAME)
        .ackPolicy(AckPolicy.Explicit)
        // ackWait — сколько сервер ждёт ack, прежде чем считать сообщение
        // недоставленным и переслать его снова (redelivery). maxDeliver —
        // сколько раз повторять, прежде чем сдаться.
        .ackWait(Duration.ofSeconds(10))
        .maxDeliver(5)
        .filterSubject(STREAM_SUBJECT)
        .build());

Дедупликация при публикации. Заголовок Nats-Msg-Id (ст. 2) в jnats задаётся не напрямую, а через messageId в PublishOptions — библиотека сама выставляет заголовок при отправке:

byte[] payload = "{\"id\":1}".getBytes(StandardCharsets.UTF_8);
PublishOptions pubOpts = PublishOptions.builder().messageId("order-1").build();

PublishAck ack1 = js.publish("orders.new", payload, pubOpts);
// ack1.getSeqno() == 1, ack1.isDuplicate() == false

PublishAck ack2 = js.publish("orders.new", payload, pubOpts);
// ack2.getSeqno() == 1, ack2.isDuplicate() == true — сервер отбросил повтор

Прогон JetStreamPull против живого nats/01-jetstream печатает это напрямую:

publish #1: seq=1, duplicate=false
publish #2 (тот же messageId): seq=1, duplicate=true

Кроме синхронного publish, у JetStream есть publishAsync, возвращающий CompletableFuture<PublishAck> — идиома, знакомая по request(): результат разбирается позже, без ожидания на каждое сообщение по очереди. Это не отдельный «второй» способ, как PublishAsync в Go с PubAckFuture и отдельными каналами Ok()/Err() — в jnats то же самое естественно ложится на CompletableFuture из стандартной библиотеки Java, без специального типа.

Pull-consumer: ConsumerContext. «Упрощённый» (simplified) API jnatsStreamContext/ConsumerContext — прямой аналог структуры Go-пакета jetstream: сначала получаем контекст consumer’а, затем читаем через него.

StreamContext streamContext = js.getStreamContext(STREAM_NAME);
ConsumerContext consumerContext = streamContext.getConsumerContext(CONSUMER_NAME);

ConsumerContext даёт, как и в Go, два режима чтения. fetchMessages(int) забирает фиксированный батч и возвращает FetchConsumerAutoCloseable с методом nextMessage(), который отдаёт очередное сообщение или null, когда батч исчерпан или истёк таймаут:

try (FetchConsumer fetchConsumer = consumerContext.fetchMessages(1)) {
    Message msg;
    while ((msg = fetchConsumer.nextMessage()) != null) {
        var meta = msg.metaData();
        System.out.printf("получено: %s: %s (попытка доставки: %d)%n",
                msg.getSubject(), new String(msg.getData(), StandardCharsets.UTF_8),
                meta.deliveredCount());
        msg.ack();
    }
}

consume(handler) — аналог Go-варианта Consume() для долгоживущего воркера: регистрирует колбэк, вызываемый на каждое новое сообщение, а управление pull-запросами под капотом библиотека берёт на себя:

MessageConsumer mc = consumerContext.consume(msg -> {
    // обработка
    msg.ack();
});
// ...
mc.stop();  // не начинать новые pull-запросы
mc.close(); // отписаться

Демо JetStreamPull использует fetchMessages, чтобы процесс завершался предсказуемо после одного сообщения — так же, как Go-демо использует Fetch по той же причине; в реальном сервисе, который живёт постоянно, consume обычно удобнее.

Ack — не единственный вариант ответа, набор методов почти буквально совпадает с Go: msg.ack() подтверждает обработку, msg.ackSync(timeout) — то же самое, но дожидается подтверждения от сервера, что ack принят (пригождается перед необратимым побочным эффектом). msg.nak()/msg.nakWithDelay(d) — явный отказ, сообщение будет доставлено снова, опционально не раньше чем через d. msg.inProgress() продлевает ackWait, не завершая обработку, — для длинных задач. msg.term() — прекратить redelivery без успеха, даже если maxDeliver не исчерпан, — для сообщений, которые в принципе невозможно обработать.

Мост от Kafka и RabbitMQ. addStream/addOrUpdateConsumer закрывают ту же задачу, что декларация топика или объявление exchange/queue — идемпотентное (для consumer) или полу-идемпотентное (для stream, с ручной обвязкой) объявление топологии при старте сервиса. fetchMessages по духу ближе всего к poll() Kafka-консьюмера: клиент явно приходит за очередной порцией и сам решает темп; 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), с собственным набором методов на интерфейсе KeyValue:

KeyValueManagement kvm = nc.keyValueManagement();
kvm.create(KeyValueConfiguration.builder().name("config").build());

KeyValue kv = nc.keyValue("config");

long rev = kv.put("feature.flag", "true".getBytes(StandardCharsets.UTF_8));

KeyValueEntry entry = kv.get("feature.flag");
// entry.getValueAsString(), entry.getRevision()

watch — то место, где идиома jnats заметнее всего расходится с Go. В новом пакете jetstream у Go watcher — канал Updates(), из которого код читает события в цикле. В jnats watch — callback-интерфейс KeyValueWatcher с двумя методами, а не одним: watch(entry) вызывается на каждое обновление, endOfData() — один раз, когда история «довычитана» и начинаются живые события. Это прямой структурный аналог nil-разделителя в Go-канале, только выражен явным методом, а не значением-маркером:

var watchSub = kv.watch("feature.flag", new KeyValueWatcher() {
    @Override
    public void watch(KeyValueEntry entry) {
        System.out.printf("watch: %s = %s (revision %d)%n",
                entry.getKey(), entry.getValueAsString(), entry.getRevision());
    }

    @Override
    public void endOfData() {
        // история довычитана, дальше только живые обновления
    }
});
// ...
watchSub.close();

KeyValueWatcher — не функциональный интерфейс (два абстрактных метода), поэтому лямбдой его не заменить, нужен анонимный класс или именованная реализация. Как и в Go, watch сразу присылает текущее значение ключа первым событием — это ожидаемо, не баг.

Замечания

Блокирующая модель, а не только async. jnats исторически построен вокруг блокирующих вызовов (request(...).get(), nextMessage(), js.publish() без Async) с отдельными неблокирующими альтернативами там, где они нужны (publishAsync, CompletableFuture из request()). Go-клиент, наоборот, по умолчанию async — обработчики подписок всегда в отдельной горутине, а блокирующее чтение (SubscribeSync/NextMsg) — осознанное исключение. На практике это означает, что в Java-коде решение «какой поток обрабатывает сообщение» видно явно — либо это поток Dispatcher, либо поток, вызвавший nextMessage()/fetchMessages(), — тогда как в Go горутины настолько дёшевы, что вопрос обычно не встаёт.

Spring и фреймворки. В этом demo jnats используется напрямую, без обёртки — так нагляднее видно, что делает сам клиент, а не абстракция поверх него. В проде, если проект уже на Spring, есть смысл посмотреть на spring-nats/аналогичные стартеры сообщества (в отличие от Kafka и RabbitMQ, у которых spring-kafka и spring-amqp — часть официальной экосистемы Spring, у NATS сопоставимого по статусу официального Spring-модуля нет) — они дают декларативные @NatsListener-подобные аннотации и управление жизненным циклом Connection через Spring-контекст, ценой ещё одного слоя абстракции над и без того простым API. Для сервиса с парой подписок разница обычно не оправдывает переход, для крупного приложения с десятками consumer’ов — вопрос стоит рассмотреть отдельно, вне рамок этой статьи.

Сопоставление идиом с Go-статьёй, если нужен быстрый справочник:

Задача nats.go jnats
Опции подключения []nats.Option (функциональные опции) Options.Builder (fluent builder)
События соединения DisconnectErrHandler/ReconnectHandler/ClosedHandler ConnectionListener.connectionEvent(conn, Events)
Ошибки протокола/подписок те же коллбэки + AsyncErrorHandler ErrorListener
Async-подписка горутина на каждую подписку Dispatcher (общий поток на несколько подписок)
Sync-подписка SubscribeSync + NextMsg subscribe(subject) + nextMessage(timeout)
Request с таймаутом RequestWithContext/Request(..., timeout) request(...).get(timeout, unit) или блокирующий request(..., timeout)
JetStream: управление jetstream.JetStream (единый объект) JetStream + JetStreamManagement (раздельно)
Идемпотентное создание stream CreateOrUpdateStream нет, только addStream/updateStream
Идемпотентное создание consumer CreateOrUpdateConsumer addOrUpdateConsumer
Dedup при публикации jetstream.WithMsgID(id) PublishOptions.builder().messageId(id)
Pull-consumer, разовое чтение consumer.Fetch(n, ...) consumerContext.fetchMessages(n)
Pull-consumer, долгоживущий consumer.Consume(handler) consumerContext.consume(handler)
KV watch канал watcher.Updates(), nil-разделитель KeyValueWatcher.watch()/endOfData()
Graceful shutdown Drain() (блокирует, таймаут через nats.DrainTimeout()) drain(timeout)CompletableFuture<Boolean>

Вывод

jnats закрывает те же задачи, что и nats.go, — reconnect с коллбэками из коробки, drain() как отдельный примитив graceful shutdown, request/reply без временных очередей, dedup на публикации через один параметр — но делает это в идиомах Java: Options.Builder вместо функциональных опций, CompletableFuture вместо каналов и context’ов, Dispatcher как явная единица разделяемого потока обработки. Упрощённый API JetStream (JetStream/JetStreamManagement/StreamContext/ConsumerContext) структурно почти буквально повторяет пакет jetstream из Go-статьи — раздельные интерфейсы для управления и потребления, fetch/consume как два явных режима pull-consumer. Единственное системное расхождение — отсутствие единого идемпотентного «создать-или-обновить» для streams, которое в demo закрыто явной обработкой кода ошибки.

Правила, вынесенные в Go-статье, не зависят от языка клиента и здесь не переписываются: at-most-once в Core, идемпотентность обработчика в JetStream (messageId защищает от дублирующей записи, но не от повторной доставки), AckWait/MaxDeliver как осознанный выбор, а не значение по умолчанию. Восьмая статья серии — security: accounts, JWT-аутентификация через nsc, mTLS; девятая — эксплуатация: наблюдаемость, backup, продакшн-чеклист.

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

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

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

Комментарии