not accepting clients
← back to blog

messaging with nats

nats routes messages by subject with a single small binary, and its jetstream layer adds persistence, replay, and at-least-once delivery on top.

services start out talking over http. one service calls another, the caller waits, and the dependency graph is whatever the call graph happens to be. this works until a third service needs the same event, until the callee is down and the caller has nowhere to put the request, or until someone needs to know what happened yesterday.

the usual next step is a broker. kafka is the default answer, and it brings partitions, consumer group rebalancing, retention tuning, and a cluster that needs its own team. nats sits at a different point on that curve: a single ~20mb binary that does subject-based pub/sub with no persistence at all, and an opt-in layer called jetstream that adds durable streams for the subjects that need them.

the short version

nats routes messages by subject. publishers send to a subject like orders.created, subscribers express interest in orders.created, orders.*, or orders.>, and the server matches them. nothing is persisted by default - a message with no interested subscriber is discarded. declaring a jetstream stream over a subject pattern changes that: matching messages are written to disk, kept under a retention policy, and delivered to consumers with acknowledgement and redelivery.

subjects and wildcards

a subject is a dot-delimited string. the tokens carry meaning, and the hierarchy is the routing structure.

orders.created
orders.eu-west.created
orders.eu-west.cancelled
metrics.node-01.cpu

subscriptions match with two wildcards. * matches exactly one token, > matches one or more trailing tokens.

orders.*.created    matches orders.eu-west.created, not orders.created
orders.>            matches every subject under orders
metrics.*.cpu       matches cpu metrics from any node

the subject convention that works is broad-to-narrow: domain, then scope, then event. putting the variable token last (orders.created.eu-west) makes it awkward to subscribe to a single region, because orders.*.eu-west also requires knowing the event name. putting it in the middle keeps both axes selectable.

subjects are not declared anywhere. there is no create-topic step in core nats. publishing to a subject that has never been used is a valid operation and costs nothing.

core pub/sub

the go client connects, subscribes, and publishes with no broker-side setup:

package main

import (
    "log"
    "time"

    "github.com/nats-io/nats.go"
)

func main() {
    nc, err := nats.Connect("nats://localhost:4222")
    if err != nil {
        log.Fatal(err)
    }
    defer nc.Drain()

    _, err = nc.Subscribe("orders.>", func(m *nats.Msg) {
        log.Printf("subject=%s payload=%s", m.Subject, m.Data)
    })
    if err != nil {
        log.Fatal(err)
    }

    nc.Publish("orders.eu-west.created", []byte(`{"id":"1042"}`))
    nc.Flush()

    time.Sleep(time.Second)
}

Publish is asynchronous. it writes to a buffer and returns; Flush waits for the server to acknowledge receipt of everything buffered so far. Drain on shutdown stops new messages, processes what’s already delivered to the subscription handlers, and then closes the connection, which is the difference between a clean deploy and dropped in-flight work.

delivery here is at-most-once and there is no memory. a subscriber that is disconnected when the message is published does not get it later. this is the correct tradeoff for telemetry, cache invalidation, and service discovery, and the wrong one for orders.

reconnection is handled by the client. it keeps a list of servers, reconnects with backoff, and resubscribes automatically:

nc, err := nats.Connect("nats://localhost:4222",
    nats.MaxReconnects(-1),
    nats.ReconnectWait(2*time.Second),
    nats.DisconnectErrHandler(func(_ *nats.Conn, err error) {
        log.Printf("disconnected: %v", err)
    }),
    nats.ReconnectHandler(func(nc *nats.Conn) {
        log.Printf("reconnected to %s", nc.ConnectedUrl())
    }),
)

messages published while disconnected are buffered in the client and flushed on reconnect, up to ReconnectBufSize. past that limit, publishes fail rather than silently vanishing.

queue groups

a plain subscription means every subscriber gets every message. for a service running three replicas, that is three times the work. a queue group makes the server pick one member per message.

nc.QueueSubscribe("orders.created", "order-processor", func(m *nats.Msg) {
    process(m.Data)
})

every replica subscribes with the same queue name. the server load-balances across the group, and a replica that disconnects is removed from rotation immediately.

                        ┌──────────────────────┐
                        │  orders.created      │
                        └──────────┬───────────┘
                                   │
              ┌────────────────────┼────────────────────┐
              ▼                    ▼                    ▼
   ┌────────────────────┐  ┌──────────────┐  ┌────────────────────┐
   │ queue: processor   │  │  audit-log   │  │ queue: notifier    │
   │  ┌────┐  ┌────┐    │  │              │  │  ┌────┐  ┌────┐    │
   │  │ p1 │  │ p2 │    │  │ plain sub -  │  │  │ n1 │  │ n2 │    │
   │  └────┘  └────┘    │  │ gets all     │  │  └────┘  └────┘    │
   │  one of the two    │  │              │  │  one of the two    │
   └────────────────────┘  └──────────────┘  └────────────────────┘

groups are independent of each other. each queue group receives the message once, and plain subscribers still receive their own copy. adding a consumer of an existing subject requires no change to the publisher, which is the property that makes subject-based routing hold up as the number of services grows.

queue membership is not configured anywhere. it exists as long as at least one client subscribes with that queue name, and disappears when the last one leaves.

request-reply

nats models request/response as pub/sub with a reply subject. the client publishes to the target subject with a temporary inbox subject attached, then waits for a message on the inbox.

the responder reads m.Reply and answers on it:

nc.QueueSubscribe("inventory.check", "inventory", func(m *nats.Msg) {
    available := checkStock(string(m.Data))
    if available {
        m.Respond([]byte(`{"in_stock":true}`))

        return
    }

    m.Respond([]byte(`{"in_stock":false}`))
})

the requester uses RequestWithContext so the timeout participates in the caller’s deadline:

ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()

msg, err := nc.RequestWithContext(ctx, "inventory.check", []byte("sku-1042"))
if err != nil {
    return fmt.Errorf("inventory check: %w", err)
}

the caller never names an instance. combining request-reply with a queue group gives load balancing and failover with no proxy, no service registry, and no client-side endpoint list. nats.ErrNoResponders comes back immediately when nothing is subscribed to the subject, which turns a dependency being down into a fast error rather than a timeout.

jetstream streams

jetstream is the persistence layer. it is enabled per server and used per subject: a stream captures a set of subject patterns and writes matching messages to a log.

streams are created from the client, so they can live in the same code that uses them:

package main

import (
    "context"
    "log"
    "time"

    "github.com/nats-io/nats.go"
    "github.com/nats-io/nats.go/jetstream"
)

func main() {
    nc, err := nats.Connect("nats://localhost:4222")
    if err != nil {
        log.Fatal(err)
    }
    defer nc.Drain()

    js, err := jetstream.New(nc)
    if err != nil {
        log.Fatal(err)
    }

    ctx := context.Background()

    stream, err := js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
        Name:       "ORDERS",
        Subjects:   []string{"orders.>"},
        Retention:  jetstream.LimitsPolicy,
        Storage:    jetstream.FileStorage,
        Replicas:   3,
        MaxAge:     7 * 24 * time.Hour,
        MaxBytes:   10 * 1024 * 1024 * 1024,
        Discard:    jetstream.DiscardOld,
        Duplicates: 2 * time.Minute,
    })
    if err != nil {
        log.Fatal(err)
    }

    log.Printf("stream %s ready", stream.CachedInfo().Config.Name)
}

Replicas: 3 puts the stream under raft. writes are acknowledged once a majority of replicas have them, and the group elects a new leader if the current one fails. MaxAge and MaxBytes are the retention bounds, and DiscardOld drops the oldest messages when a bound is hit rather than rejecting new writes.

publishing through jetstream returns an acknowledgement with the assigned sequence number:

ack, err := js.Publish(ctx, "orders.created", payload,
    jetstream.WithMsgID("order-1042-created"))
if err != nil {
    return fmt.Errorf("publish: %w", err)
}

log.Printf("stored at %s seq %d duplicate=%t", ack.Stream, ack.Sequence, ack.Duplicate)

WithMsgID combined with the stream’s Duplicates window gives idempotent publishing. a retried publish with the same message id inside the window is not stored twice, and ack.Duplicate reports that it was suppressed. this is what makes retrying a failed publish safe, and it is the cheapest available fix for double-writes on a flaky network.

the same message is now both live pub/sub and a durable record. an existing core subscriber on orders.created is unaffected by the stream being added.

consumers and acknowledgement

a stream stores messages. a consumer tracks a position in it. multiple consumers read the same stream independently, each with its own cursor, filter, and ack state.

cons, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
    Durable:       "fulfilment",
    FilterSubject: "orders.created",
    AckPolicy:     jetstream.AckExplicitPolicy,
    AckWait:       30 * time.Second,
    MaxDeliver:    5,
    MaxAckPending: 256,
    DeliverPolicy: jetstream.DeliverAllPolicy,
})
if err != nil {
    log.Fatal(err)
}

Durable names the consumer so its position survives restarts and disconnects. MaxAckPending is the flow control: the server will not have more than 256 unacknowledged messages outstanding, which bounds how far ahead of the handler the delivery can run.

consuming is a callback over a long-lived pull subscription:

cc, err := cons.Consume(func(msg jetstream.Msg) {
    if err := handle(msg.Data()); err != nil {
        msg.NakWithDelay(5 * time.Second)

        return
    }

    msg.Ack()
})
if err != nil {
    log.Fatal(err)
}
defer cc.Stop()

the four terminal outcomes are the whole delivery contract. Ack marks the message done and advances the cursor. Nak asks for immediate redelivery, and NakWithDelay backs off first. Term refuses the message permanently, so it is never redelivered regardless of MaxDeliver. doing nothing means redelivery after AckWait expires.

for work that legitimately takes longer than AckWait, msg.InProgress() resets the timer:

cc, err := cons.Consume(func(msg jetstream.Msg) {
    done := make(chan error, 1)
    go func() { done <- transcode(msg.Data()) }()

    for {
        select {
        case err := <-done:
            if err != nil {
                msg.Nak()

                return
            }
            msg.Ack()

            return
        case <-time.After(15 * time.Second):
            msg.InProgress()
        }
    }
})

adding a second replica of the same service with the same Durable name gives the queue-group behaviour with persistence: the server distributes messages across both, and unacknowledged work from a replica that dies is redelivered to the other. unlike kafka consumer groups, this involves no partition assignment: two replicas of a one-subject consumer both make progress.

when the handler has exhausted MaxDeliver, the message is dropped unless the stream has a dead-letter path. jetstream publishes an advisory for it:

nc.Subscribe("$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.ORDERS.fulfilment",
    func(m *nats.Msg) {
        log.Printf("giving up on message: %s", m.Data)
    })

subscribing to that advisory and republishing the referenced message to a dlq.> subject is the standard pattern. it is not built in.

the key-value store

jetstream streams have enough machinery to build a key-value store, and the client exposes one directly. it is a stream with LimitsPolicy retention, subject-per-key, and history depth as MaxMsgsPerSubject.

kv, err := js.CreateKeyValue(ctx, jetstream.KeyValueConfig{
    Bucket:   "feature-flags",
    History:  5,
    Storage:  jetstream.FileStorage,
    Replicas: 3,
})
if err != nil {
    log.Fatal(err)
}

rev, err := kv.Put(ctx, "checkout.new-flow", []byte("enabled"))

reads return the value with its revision, and updates can be conditional on it:

entry, err := kv.Get(ctx, "checkout.new-flow")
if err != nil {
    return err
}

_, err = kv.Update(ctx, "checkout.new-flow", []byte("disabled"), entry.Revision())
if err != nil {
    return fmt.Errorf("lost the race: %w", err)
}

Update fails if the revision has moved, which is compare-and-swap. that is enough to build a leader election or a distributed lock without another dependency.

a watcher streams changes as they happen, replaying current values first:

w, err := kv.WatchAll(ctx)
if err != nil {
    return err
}
defer w.Stop()

for entry := range w.Updates() {
    if entry == nil {
        log.Print("initial state replayed")

        continue
    }

    log.Printf("%s = %s (rev %d, %s)", entry.Key(), entry.Value(),
        entry.Revision(), entry.Operation())
}

the single nil on the channel marks the end of the initial replay and the start of live updates. that boundary is what makes it usable for configuration: block until nil, start serving with a complete picture, then apply changes as they arrive. entry.Operation() distinguishes a KeyValuePut from a KeyValueDelete or a KeyValuePurge.

running on kubernetes

the official helm chart runs nats as a statefulset with jetstream on persistent volumes:

helm repo add nats https://nats-io.github.io/k8s/helm/charts/

a three-node clustered deployment with file-backed jetstream:

# values.yaml
config:
  cluster:
    enabled: true
    replicas: 3
  jetstream:
    enabled: true
    fileStore:
      pvc:
        size: 50Gi
        storageClassName: gp3
  merge:
    max_payload: 4MB

podTemplate:
  topologySpreadConstraints:
    kubernetes.io/hostname:
      maxSkew: 1
      whenUnsatisfiable: DoNotSchedule

promExporter:
  enabled: true
  podMonitor:
    enabled: true

install it into its own namespace:

helm upgrade --install nats nats/nats \
  --namespace nats \
  --create-namespace \
  --values values.yaml

three replicas is the minimum for R3 streams, since raft needs a majority to accept writes. the topology spread constraint keeps the three pods on separate nodes, without which a single node failure can take down two thirds of the quorum.

the storage class matters for the same reason it matters for postgres. jetstream file storage does synchronous writes on the raft path, and volumes with high tail latency show up directly as publish latency.

clients inside the cluster connect to the headless service:

nc, err := nats.Connect("nats://nats.nats.svc.cluster.local:4222")

the client discovers the rest of the cluster from the server’s INFO message on connect, so a single dns name is enough to get failover across all three pods.

operating it

the nats cli is the primary operational tool. it talks the same protocol as any other client and needs no server-side agent.

nats --server nats://localhost:4222 stream report

this lists every stream with message counts, byte size, consumer count, and the current raft leader. for a single consumer’s health:

nats consumer info ORDERS fulfilment

the two numbers to read there are unprocessed messages and redelivered. a growing unprocessed count means the consumer cannot keep up with the publish rate. a growing redelivery count means handlers are failing or exceeding AckWait, and the two have different fixes.

the monitoring endpoint on port 8222 exposes the same data as json, and the prometheus exporter enabled in the chart above turns it into metrics. the ones worth alerting on:

  • jetstream_consumer_num_pending - consumer lag, the closest analogue to kafka’s consumer lag
  • jetstream_consumer_num_redelivered - handler failures and ack timeouts
  • gnatsd_varz_slow_consumers - clients the server disconnected for not keeping up, and should stay at zero
  • jetstream_server_total_streams and store bytes against the configured limits
  • raft leader changes on jetstream streams, which indicate node or disk instability

for capacity work, the cli also benchmarks against a real cluster:

nats bench benchsubject --pub 4 --sub 4 --msgs 500000 --size 1024

running this against the actual storage class and node types gives a number that means something, which is not true of the throughput figures in any comparison blog post.

what to watch out for

core nats forgets everything. a message with no connected subscriber is discarded, silently and by design. this is the single most common surprise when moving from a queue-based broker. anything that must survive a consumer restart needs a jetstream stream, and adding one later does not backfill the messages already lost.

streams cannot have overlapping subjects. creating a stream whose subject pattern overlaps an existing stream’s is rejected. a stream on orders.> blocks a later stream on orders.created, which forces the subject hierarchy to be planned before the second team needs a stream. splitting an over-broad stream after the fact means creating a new one and replaying.

handlers must be idempotent. jetstream is at-least-once. an ack lost on the network, a handler that overruns AckWait, or a consumer that crashes after the side effect and before the ack all produce a redelivery. deduplicating on msg.Headers().Get(jetstream.MsgIDHeader) or on a natural key in the payload is not optional.

MaxAckPending is the throughput limit. the default is 1000. a consumer whose handler takes 100ms per message and never raises the limit is capped regardless of how many replicas run, because the server stops delivering past the outstanding window. it is the first thing to change when a consumer is inexplicably slow.

slow core subscribers drop messages. each subscription has a client-side pending limit, 500,000 messages or 64mb by default. a handler slower than the publish rate fills it, and the client starts dropping messages for that subscription and reports nats.ErrSlowConsumer to the async error handler. separately, the server disconnects connections that fall too far behind, which is what gnatsd_varz_slow_consumers counts. registering nats.ErrorHandler at connect time is the difference between noticing either case and not.

payloads are small by default. max_payload is 1mb, and raising it much past a few megabytes hurts everything sharing the connection. large payloads belong in the jetstream object store or object storage, with the message carrying a reference.

authentication is not on by default. an unconfigured server accepts any connection on 4222. production deployments want either nkeys or the operator/account model with jwts, which is also what provides multi-tenancy through account isolation. deciding this after clients exist means rotating every client’s credentials at once.

references

[1] nats documentation. “subject-based messaging.”
docs.nats.io/nats-concepts/subjects

[2] nats documentation. “queue groups.”
docs.nats.io/nats-concepts/core-nats/queue

[3] nats documentation. “jetstream.”
docs.nats.io/nats-concepts/jetstream

[4] nats documentation. “consumers.”
docs.nats.io/nats-concepts/jetstream/consumers

[5] nats documentation. “key/value store.”
docs.nats.io/nats-concepts/jetstream/key-value-store

[6] nats documentation. “monitoring.”
docs.nats.io/running-a-nats-service/nats_admin/monitoring

[7] nats-io. “nats.go - the go client.”
github.com/nats-io/nats.go

[8] nats-io. “nats helm charts.”
github.com/nats-io/k8s

# ask the author

a question
about this
post?

direct line