Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

NATS

A high-performance, cloud-native messaging system designed for modern distributed systems with support for pub/sub, request/reply, and streaming.

NATS Cheatsheet

A high-performance, cloud-native messaging system designed for modern distributed systems with support for pub/sub, request/reply, and streaming.

Overview

NATS is a simple, secure, and performant communications system for digital systems, services, and devices. Originally built by Derek Collison in 2010, NATS is now a Cloud Native Computing Foundation (CNCF) project. It provides core messaging (pub/sub), distributed queueing (queue groups), request/reply, and persistent streaming via JetStream.

NATS ArchitectureJetStreamSubscribersPublishersNATS Server ClusterClusterClusterClusterPublishPublishSubscribeSubscribeSubscribePersistPersistPersistStream StorageNATS Server 1NATS Server 2NATS Server 3Service AService BWorker 1Worker 2DashboardNATS ArchitectureJetStreamSubscribersPublishersNATS Server ClusterClusterClusterClusterPublishPublishSubscribeSubscribeSubscribePersistPersistPersistStream StorageNATS Server 1NATS Server 2NATS Server 3Service AService BWorker 1Worker 2Dashboard

Core Concepts

Subjects and Wildcards

Key Concepts:

  • Subject: A UTF-8 string that messages are published to and subscribed from (e.g., orders.created, sensors.temperature.london)
  • Subject Hierarchy: Dot-separated tokens create a hierarchy (e.g., app.service.region.instance)
  • Wildcards: Pattern matching for subscriptions
    • * - Matches exactly one token (e.g., sensors.*.london matches sensors.temperature.london)
    • > - Matches one or more tokens at the end (e.g., orders.> matches orders.created, orders.updated.paid)
  • Subject-Based Routing: Intelligent routing without broker configuration

Common Patterns

Subject Naming Conventions:

# Entity-based subjects
users.created
users.updated
users.deleted

# Hierarchical subjects
region.eu.service.payment.event.transaction
app.microservice.log.error

# Time-series data
metrics.cpu.host1.2024.01.15
sensors.temperature.building-a.floor-2.room-5

# Request-reply subjects
service.api.user.get
service.api.order.create

Wildcard Examples:

# Subscribe to all temperature sensors
nats sub "sensors.temperature.*"

# Subscribe to all events for a specific order
nats sub "orders.12345.>"

# Subscribe to all error logs across services
nats sub "*.*.log.error"

# Subscribe to all metrics for a specific host
nats sub "metrics.>.host1"

Subject Design Best Practices

# Good: Clear hierarchy, scalable
payments.transactions.created
payments.transactions.approved
payments.refunds.processed

# Avoid: Too generic
data
events
messages

# Good: Include context
eu-west-1.inventory.stock.low
us-east-1.inventory.stock.low

# Avoid: Hard to filter
inventory-stock-low-eu-west-1

Publish/Subscribe and Queue Groups

Basic Pub/Sub

Key Concepts:

  • Publisher: Client that sends messages to a subject
  • Subscriber: Client that receives messages from a subject
  • Fan-Out: One message published to multiple subscribers
  • Fire-and-Forget: No acknowledgement by default (at-most-once delivery)
Publish to'orders.created'DeliverDeliverDeliverPublisherNATS ServerSubscriber 1Subscriber 2Subscriber 3Publish to'orders.created'DeliverDeliverDeliverPublisherNATS ServerSubscriber 1Subscriber 2Subscriber 3

Publishing Messages:

# Simple publish
nats pub orders.created "Order #123 created"

# Publish JSON payload
nats pub users.created '{"id": "user-456", "email": "user@example.com"}'

# Publish with headers
nats pub --header "Content-Type:application/json" events.data '{"value": 42}'

# Request-reply pattern (wait for response)
nats req service.time.get ""

Subscribing to Messages:

# Simple subscription
nats sub orders.created

# Subscribe with wildcards
nats sub "orders.>"

# Subscribe (headers are rendered automatically when present)
nats sub orders.created

# Show only headers, suppressing the message body
nats sub --headers-only orders.created

# Subscribe with queue group
nats sub --queue workers orders.process

Queue Groups

Key Concepts:

  • Queue Group: Load balancing pattern where only one member receives each message
  • Queue Subscriber: Member of a queue group
  • Random Distribution: Messages distributed randomly among queue members
  • Horizontal Scaling: Add more queue members to scale processing
Queue Group 'workers'orders.processMsg 1Msg 2Msg 3Msg 4PublisherNATS ServerWorker 1Worker 2Worker 3Queue Group 'workers'orders.processMsg 1Msg 2Msg 3Msg 4PublisherNATS ServerWorker 1Worker 2Worker 3

Queue Group Examples:

# Start multiple workers in same queue group
# Terminal 1
nats sub --queue processors orders.process

# Terminal 2
nats sub --queue processors orders.process

# Terminal 3
nats sub --queue processors orders.process

# Publish messages - only one worker receives each message
nats pub orders.process "Process order 1"
nats pub orders.process "Process order 2"
nats pub orders.process "Process order 3"

Code Examples

Go - Publishing and Subscribing:

package main

import (
    "fmt"
    "log"
    "time"

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

func main() {
    // Connect to NATS
    nc, err := nats.Connect(nats.DefaultURL)
    if err != nil {
        log.Fatal(err)
    }
    defer nc.Close()

    // Simple publish
    err = nc.Publish("orders.created", []byte("Order #123"))
    if err != nil {
        log.Fatal(err)
    }

    // Subscribe
    sub, err := nc.Subscribe("orders.created", func(m *nats.Msg) {
        fmt.Printf("Received: %s\n", string(m.Data))
    })
    if err != nil {
        log.Fatal(err)
    }
    defer sub.Unsubscribe()

    // Queue group subscription
    nc.QueueSubscribe("orders.process", "workers", func(m *nats.Msg) {
        fmt.Printf("Worker processing: %s\n", string(m.Data))
    })

    // Request-reply pattern
    msg, err := nc.Request("service.time", []byte(""), 1*time.Second)
    if err != nil {
        log.Fatal(err)
    }
    fmt.Printf("Reply: %s\n", string(msg.Data))

    // Keep alive
    time.Sleep(10 * time.Second)
}

Python - Using nats.py:

import asyncio
import nats
from nats.errors import TimeoutError

async def main():
    # Connect to NATS
    nc = await nats.connect("nats://localhost:4222")

    # Simple publish
    await nc.publish("orders.created", b"Order #123")

    # Subscribe
    async def message_handler(msg):
        print(f"Received: {msg.data.decode()}")

    await nc.subscribe("orders.created", cb=message_handler)

    # Queue group subscription
    async def worker_handler(msg):
        print(f"Worker processing: {msg.data.decode()}")

    await nc.subscribe("orders.process", queue="workers", cb=worker_handler)

    # Request-reply
    try:
        response = await nc.request("service.time", b"", timeout=1.0)
        print(f"Reply: {response.data.decode()}")
    except TimeoutError:
        print("Request timeout")

    # Keep alive
    await asyncio.sleep(10)
    await nc.close()

if __name__ == '__main__':
    asyncio.run(main())

JavaScript/Node.js - Using nats.js:

import { connect, StringCodec } from 'nats';

async function main() {
    // Connect to NATS
    const nc = await connect({ servers: 'nats://localhost:4222' });
    const sc = StringCodec();

    // Simple publish
    nc.publish('orders.created', sc.encode('Order #123'));

    // Subscribe
    const sub = nc.subscribe('orders.created');
    (async () => {
        for await (const m of sub) {
            console.log(`Received: ${sc.decode(m.data)}`);
        }
    })();

    // Queue group subscription
    const qsub = nc.subscribe('orders.process', { queue: 'workers' });
    (async () => {
        for await (const m of qsub) {
            console.log(`Worker processing: ${sc.decode(m.data)}`);
        }
    })();

    // Request-reply
    const response = await nc.request('service.time', sc.encode(''), { timeout: 1000 });
    console.log(`Reply: ${sc.decode(response.data)}`);

    // Cleanup
    setTimeout(() => {
        nc.drain();
    }, 10000);
}

main();

JetStream Basics

Overview

JetStream is NATS 2.0's built-in persistence layer providing:

  • Stream storage (at-least-once, exactly-once semantics)
  • Message replay
  • Acknowledgements
  • Durable consumers
  • Message deduplication
JetStream ArchitecturePublishPull/PushPull/PushPull/PushAckAckAckRetention PolicyPublisherStreamorders.createdSubjects: orders.*Consumer 1DurableConsumer 2EphemeralConsumer 3Queue GroupRetentionLimits/Age/InterestJetStream ArchitecturePublishPull/PushPull/PushPull/PushAckAckAckRetention PolicyPublisherStreamorders.createdSubjects: orders.*Consumer 1DurableConsumer 2EphemeralConsumer 3Queue GroupRetentionLimits/Age/Interest

Streams

Key Concepts:

  • Stream: Persistent message store bound to one or more subjects
  • Retention Policy: How messages are deleted (limits, age, interest-based)
  • Storage: File or memory-based
  • Replication: Cluster replication factor
  • Limits: Max messages, bytes, age, message size

Creating Streams:

# Create a stream with nats CLI
nats stream add orders \
    --subjects "orders.*" \
    --retention limits \
    --storage file \
    --replicas 3 \
    --max-msgs 1000000 \
    --max-bytes 1GB \
    --max-age 7d \
    --max-msg-size 1MB \
    --discard old

# Create stream (interactive mode)
nats stream add

# List streams
nats stream ls

# View stream info
nats stream info orders

# Edit stream configuration
nats stream edit orders

# Delete stream
nats stream rm orders

# Purge all messages from stream
nats stream purge orders

# View stream messages
nats stream view orders

Stream Configuration Examples:

# Stream for logs (retention by age)
nats stream add logs \
    --subjects "logs.>" \
    --retention limits \
    --max-age 30d \
    --storage file

# Stream for metrics (retention by size)
nats stream add metrics \
    --subjects "metrics.>" \
    --retention limits \
    --max-bytes 10GB \
    --discard old

# Stream for events (work queue - message removed once acked)
nats stream add events \
    --subjects "events.>" \
    --retention workq \
    --storage file

# In-memory stream for fast processing
nats stream add cache \
    --subjects "cache.>" \
    --retention limits \
    --storage memory \
    --max-age 1h

Consumers

Key Concepts:

  • Consumer: View into a stream with its own delivery and acknowledgement semantics
  • Durable Consumer: Survives server restarts, remembers position
  • Ephemeral Consumer: Exists only while connected
  • Push Consumer: Server pushes messages to subscriber
  • Pull Consumer: Client pulls messages on demand
  • Ack Policy: None, All, Explicit
  • Replay Policy: Instant, Original (time-based)
Message arrivesConsumer fetchesConsumer acksAck timeoutRetryToo many retriesMove to dead letterPendingDeliveredAcknowledgedRedeliveredMaxRedeliveriesMessage arrivesConsumer fetchesConsumer acksAck timeoutRetryToo many retriesMove to dead letterPendingDeliveredAcknowledgedRedeliveredMaxRedeliveries

Creating Consumers:

# Create pull consumer (durable)
nats consumer add orders order-processor \
    --pull \
    --deliver all \
    --ack explicit \
    --max-deliver 5 \
    --max-pending 100 \
    --wait 30s

# Create push consumer (push mode is implied by --target, the delivery subject)
nats consumer add orders email-notifier \
    --target "notifications.email" \
    --deliver all \
    --ack explicit

# Create ephemeral consumer
nats consumer add orders temp-consumer \
    --ephemeral \
    --pull \
    --deliver new

# List consumers for a stream
nats consumer ls orders

# View consumer info
nats consumer info orders order-processor

# Delete consumer
nats consumer rm orders order-processor

Consuming Messages:

# Pull messages from consumer
nats consumer next orders order-processor

# Pull multiple messages
nats consumer next orders order-processor --count 10

# Subscribe to stream (creates ephemeral consumer)
nats subscribe orders.created --stream orders

# Subscribe with acknowledgements
nats subscribe orders.created --stream orders --ack

Code Examples

Go - JetStream Publishing and Consuming:

package main

import (
    "fmt"
    "log"
    "time"

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

func main() {
    // Connect
    nc, _ := nats.Connect(nats.DefaultURL)
    defer nc.Close()

    // Create JetStream context
    js, err := nc.JetStream()
    if err != nil {
        log.Fatal(err)
    }

    // Create stream
    js.AddStream(&nats.StreamConfig{
        Name:     "orders",
        Subjects: []string{"orders.*"},
        Storage:  nats.FileStorage,
        MaxAge:   7 * 24 * time.Hour,
    })

    // Publish to stream
    ack, err := js.Publish("orders.created", []byte("Order #123"))
    if err != nil {
        log.Fatal(err)
    }
    fmt.Printf("Published message, sequence: %d\n", ack.Sequence)

    // Create pull consumer
    js.AddConsumer("orders", &nats.ConsumerConfig{
        Durable:   "processor",
        AckPolicy: nats.AckExplicitPolicy,
    })

    // Subscribe with pull consumer
    sub, _ := js.PullSubscribe("orders.*", "processor")

    // Fetch messages
    msgs, _ := sub.Fetch(10, nats.MaxWait(5*time.Second))
    for _, msg := range msgs {
        fmt.Printf("Received: %s\n", string(msg.Data))
        msg.Ack()
    }

    // Push consumer
    js.Subscribe("orders.*", func(msg *nats.Msg) {
        fmt.Printf("Push received: %s\n", string(msg.Data))
        msg.Ack()
    }, nats.Durable("email-notifier"))
}

Python - JetStream:

import asyncio
import nats
from nats.js.api import StreamConfig, ConsumerConfig, AckPolicy

async def main():
    nc = await nats.connect("nats://localhost:4222")
    js = nc.jetstream()

    # Create stream
    try:
        await js.add_stream(
            name="orders",
            subjects=["orders.*"],
            max_age=7 * 24 * 60 * 60  # 7 days in seconds
        )
    except Exception as e:
        print(f"Stream might already exist: {e}")

    # Publish
    ack = await js.publish("orders.created", b"Order #123")
    print(f"Published sequence: {ack.seq}")

    # Pull consumer
    await js.add_consumer(
        stream="orders",
        config=ConsumerConfig(
            durable_name="processor",
            ack_policy=AckPolicy.EXPLICIT
        )
    )

    # Fetch and process messages
    psub = await js.pull_subscribe("orders.*", "processor")
    msgs = await psub.fetch(10, timeout=5)
    for msg in msgs:
        print(f"Received: {msg.data.decode()}")
        await msg.ack()

    await nc.close()

if __name__ == '__main__':
    asyncio.run(main())

JavaScript - JetStream:

import { connect, StringCodec, AckPolicy } from 'nats';

async function main() {
    const nc = await connect({ servers: 'nats://localhost:4222' });
    const js = nc.jetstream();
    const sc = StringCodec();

    // Create stream
    await js.streams.add({
        name: 'orders',
        subjects: ['orders.*'],
        max_age: 7 * 24 * 60 * 60 * 1000000000, // 7 days in nanoseconds
    });

    // Publish
    const ack = await js.publish('orders.created', sc.encode('Order #123'));
    console.log(`Published sequence: ${ack.seq}`);

    // Create pull consumer
    await js.consumers.add('orders', {
        durable_name: 'processor',
        ack_policy: AckPolicy.Explicit,
    });

    // Pull messages
    const consumer = await js.consumers.get('orders', 'processor');
    const iter = await consumer.fetch({ max_messages: 10, expires: 5000 });

    for await (const msg of iter) {
        console.log(`Received: ${sc.decode(msg.data)}`);
        msg.ack();
    }

    await nc.drain();
}

main();

Message Deduplication

# Enable deduplication with message ID header
nats pub orders.created \
    --header "Nats-Msg-Id:unique-msg-123" \
    "Order data"

# Duplicate publish (will be ignored within deduplication window)
nats pub orders.created \
    --header "Nats-Msg-Id:unique-msg-123" \
    "Order data"

Go - Message Deduplication:

// Publish with message ID for deduplication
js.Publish("orders.created", []byte("Order #123"), nats.MsgId("unique-msg-123"))

Connection Tuning and Clustering

Connection Configuration

Key Concepts:

  • Ping Interval: Client-server heartbeat frequency
  • Max Reconnect Attempts: Connection retry limit
  • Reconnect Wait: Delay between reconnection attempts
  • Drain: Graceful shutdown with pending message processing
  • TLS: Encrypted connections
  • Token/Credentials: Authentication methods

Connection Options:

# Connect with authentication
nats --user admin --password secret sub test

# Connect with token
nats --token mytoken123 sub test

# Connect with credentials file
nats --creds /path/to/user.creds sub test

# Connect with TLS
nats --tlscert /path/to/client-cert.pem \
     --tlskey /path/to/client-key.pem \
     --tlsca /path/to/ca.pem \
     sub test

# Connect to multiple servers (failover)
nats --server nats://server1:4222,nats://server2:4222 sub test

Go - Connection Tuning:

package main

import (
    "time"
    "github.com/nats-io/nats.go"
)

func main() {
    // Advanced connection options
    nc, _ := nats.Connect("nats://localhost:4222",
        // Reconnection
        nats.MaxReconnects(10),
        nats.ReconnectWait(2*time.Second),
        nats.ReconnectBufSize(5*1024*1024),

        // Ping/Pong
        nats.PingInterval(20*time.Second),
        nats.MaxPingsOutstanding(5),

        // Timeout
        nats.Timeout(5*time.Second),

        // Callbacks
        nats.DisconnectErrHandler(func(nc *nats.Conn, err error) {
            log.Printf("Disconnected: %v", err)
        }),
        nats.ReconnectHandler(func(nc *nats.Conn) {
            log.Printf("Reconnected to %s", nc.ConnectedUrl())
        }),
        nats.ClosedHandler(func(nc *nats.Conn) {
            log.Printf("Connection closed")
        }),

        // Multiple servers
        nats.Name("my-service"),
        nats.UserInfo("user", "password"),
    )
    defer nc.Drain() // Graceful shutdown
}

Python - Connection Configuration:

import asyncio
import nats

async def main():
    async def disconnected_cb():
        print("Disconnected from NATS")

    async def reconnected_cb():
        print("Reconnected to NATS")

    async def error_cb(e):
        print(f"Error: {e}")

    nc = await nats.connect(
        servers=["nats://server1:4222", "nats://server2:4222"],
        reconnect_time_wait=2,
        max_reconnect_attempts=10,
        ping_interval=20,
        max_outstanding_pings=5,
        user="admin",
        password="secret",
        disconnected_cb=disconnected_cb,
        reconnected_cb=reconnected_cb,
        error_cb=error_cb,
    )

    # Graceful shutdown
    await nc.drain()

asyncio.run(main())

Clustering

Key Concepts:

  • Cluster: Multiple NATS servers forming a full mesh
  • Route: Connection between cluster members
  • Gateway: Connection between separate clusters (super-cluster)
  • Leaf Node: Remote server connecting to a cluster
  • Full Mesh: Every server connects to every other server
Cluster 2 - USCluster 1 - EURouteRouteRouteRouteGatewayLeafNATS EU-1NATS EU-2NATS EU-3NATS US-1NATS US-2Leaf NodeCluster 2 - USCluster 1 - EURouteRouteRouteRouteGatewayLeafNATS EU-1NATS EU-2NATS EU-3NATS US-1NATS US-2Leaf Node

Basic Cluster Configuration:

# nats-server-1.conf
port: 4222
server_name: nats-1

cluster {
    name: main-cluster
    listen: 0.0.0.0:6222
    routes: [
        nats://nats-2:6222
        nats://nats-3:6222
    ]
}

jetstream {
    store_dir: /data/jetstream
    max_memory_store: 1GB
    max_file_store: 10GB
}
# nats-server-2.conf
port: 4222
server_name: nats-2

cluster {
    name: main-cluster
    listen: 0.0.0.0:6222
    routes: [
        nats://nats-1:6222
        nats://nats-3:6222
    ]
}

jetstream {
    store_dir: /data/jetstream
    max_memory_store: 1GB
    max_file_store: 10GB
}

Starting a Cluster:

# Start first server
nats-server -c nats-server-1.conf

# Start second server
nats-server -c nats-server-2.conf

# Start third server
nats-server -c nats-server-3.conf

# Check cluster status
nats server list
nats server info
nats server report jetstream

Gateway Configuration (Super-Cluster):

# EU cluster gateway config
gateway {
    name: eu-cluster
    listen: 0.0.0.0:7222
    gateways: [
        {name: "us-cluster", urls: ["nats://us-gateway:7222"]}
        {name: "asia-cluster", urls: ["nats://asia-gateway:7222"]}
    ]
}

Leaf Node Configuration:

# Leaf node connecting to cluster
port: 4222

leafnodes {
    remotes = [
        {
            url: "nats://cluster-server:7422"
            credentials: "/path/to/leaf.creds"
        }
    ]
}

Monitoring and Health

# Server monitoring endpoint (HTTP)
curl http://localhost:8222/varz        # Server stats
curl http://localhost:8222/connz       # Connection info
curl http://localhost:8222/routez      # Route info
curl http://localhost:8222/subsz       # Subscription info
curl http://localhost:8222/jsz         # JetStream info
curl http://localhost:8222/healthz     # Health check

# Enable monitoring in config
http_port: 8222

# nats CLI monitoring
nats server list
nats server info nats-1
nats server report connections
nats server report jetstream
nats server report accounts

Common Tooling (nats CLI)

Installation

# Install via script (Linux/macOS)
curl -sf https://binaries.nats.dev/nats-io/natscli/nats@latest | sh

# Install via Homebrew (macOS)
brew install nats-io/nats-tools/nats

# Install via Docker
docker pull natsio/nats-box:latest
docker run --rm -it natsio/nats-box:latest

# Download binary from GitHub
# https://github.com/nats-io/natscli/releases

Context Management

# Add a context (server connection profile)
nats context add production \
    --server nats://prod.example.com:4222 \
    --description "Production NATS cluster" \
    --creds /path/to/prod.creds

nats context add local \
    --server nats://localhost:4222 \
    --description "Local development"

# List contexts
nats context ls

# Select context
nats context select production

# Show current context
nats context info

# Remove context
nats context rm local

Essential Commands

Server and Connection:

# Check server connectivity
nats server check connection

# Round-trip latency test
nats rtt

# Server information (requires the system account; see note below)
nats server info

# Account information
nats account info

System account required: nats server info, nats server list, nats server check <...> (beyond connection), and nats server report ... issue requests on the $SYS system account and fail with "ensure the account used has system privileges" against a default server. Connect with the configured system-account user (e.g. via a context whose creds map to $SYS), or enable a system_account in the server config. nats server check connection, nats rtt, and nats account info work without it.

Publishing and Subscribing:

# Basic pub/sub
nats pub subject.name "message"
nats sub subject.name

# Wildcards
nats sub "sensors.>"
nats sub "metrics.*.cpu"

# Request/reply
nats req service.time ""
nats reply service.time "{{.Time}}"

# Bench testing (subcommand per role: pub, sub, js, kv, service)
nats bench pub test --clients 10 --msgs 10000
nats bench sub test --clients 5 --msgs 10000

Stream Management:

# Stream operations
nats stream ls
nats stream add STREAM_NAME
nats stream info STREAM_NAME
nats stream edit STREAM_NAME
nats stream rm STREAM_NAME
nats stream purge STREAM_NAME

# View stream data
nats stream view STREAM_NAME
nats stream get STREAM_NAME 100  # Get message by sequence

# Stream backup/restore
nats stream backup STREAM_NAME backup.tar.gz
nats stream restore STREAM_NAME backup.tar.gz

Consumer Management:

# Consumer operations
nats consumer ls STREAM_NAME
nats consumer add STREAM_NAME CONSUMER_NAME
nats consumer info STREAM_NAME CONSUMER_NAME
nats consumer rm STREAM_NAME CONSUMER_NAME

# Consume messages
nats consumer next STREAM_NAME CONSUMER_NAME
nats consumer next STREAM_NAME CONSUMER_NAME --count 10
nats consumer sub STREAM_NAME CONSUMER_NAME

Schema Tools:

nats schema works against NATS's built-in catalogue of server API/event schemas (advisories, JetStream API, etc.) — it is not a user-facing schema registry, and there is no nats schema add or nats pub --schema.

# Search the built-in schema catalogue
nats schema search jetstream

# Show a schema's contents
nats schema info io.nats.jetstream.api.v1.stream_create_request

# Validate a JSON file against a known schema
nats schema validate io.nats.jetstream.api.v1.stream_create_request msg.json

Traffic Monitoring:

# Live subscription monitoring
nats sub ">" --translate="{{.Subject}}"

# Show advisories and events (server advisories shown by default)
nats events

# Show all event types (client connections, JetStream advisories, etc.)
nats events --all

Advanced Features

Message Tracing:

# Enable message tracing
nats pub --trace orders.created "Order data"

# Results show message path through cluster

Key-Value Store:

# Create KV bucket
nats kv add config

# Put value
nats kv put config app.version "1.2.3"

# Get value
nats kv get config app.version

# Delete key
nats kv del config app.version

# List keys
nats kv ls config

# Watch for changes
nats kv watch config

Object Store:

# Create object store
nats object add files

# Put file (object name defaults to the file name; override with --name)
nats object put files /path/to/file.pdf
nats object put files /path/to/file.pdf --name myfile.pdf

# Get file (-O / --output overrides the local output path)
nats object get files myfile.pdf -O /path/to/download.pdf

# List objects
nats object ls files

# Delete object
nats object rm files myfile.pdf

Quick Reference

Subject Wildcards

Wildcard Matches Example
* Exactly one token sensors.*.room1 matches sensors.temp.room1
> One or more tokens (end only) orders.> matches orders.created.paid

Message Delivery Guarantees

Type Guarantee Use Case
Core NATS At-most-once Fast, fire-and-forget messaging
JetStream (Ack) At-least-once Reliable message delivery
JetStream (Dedup) Exactly-once Critical transactions

Consumer Types

Type Persistence Delivery Use Case
Durable Pull Survives restarts Client-controlled Worker pools, batch processing
Durable Push Survives restarts Server-controlled Real-time processing
Ephemeral Session-only Client/Server Temporary monitoring

Retention Policies

Policy Behaviour Use Case
Limits Keep until max size/age/count Logs, metrics
Interest Keep until all consumers acknowledge Work queues
WorkQueue Remove after ack (one consumer) Task distribution

Common nats CLI Commands

Command Description
nats pub SUBJECT MSG Publish message
nats sub SUBJECT Subscribe to subject
nats req SUBJECT MSG Request-reply
nats stream add Create stream
nats stream ls List streams
nats consumer add Create consumer
nats consumer next Pull next message
nats server info Server information
nats bench Performance benchmark

Connection URLs

Format Example Description
Plain nats://localhost:4222 Unencrypted connection
TLS tls://localhost:4222 TLS encrypted
Credentials nats://user:pass@localhost:4222 With authentication
WebSocket ws://localhost:8080 WebSocket connection

Common Issues and Solutions

Connection Problems

Issue Symptoms Solution
Connection refused Client cannot connect Check server is running: nats server check connection
Authentication failed 'Authorization Violation' error Verify credentials: nats --creds /path/to/file.creds sub test
Slow subscriber Messages backing up Increase processing speed or add queue group members
Connection dropped Periodic disconnects Check network stability, increase ping interval

Debugging Connection Issues:

# Test basic connectivity
nats server check connection

# Measure round-trip time
nats rtt

# Verbose connection logging
nats --trace sub test

# Check server logs
nats server report connections

JetStream Issues

Issue Cause Solution
Stream not found JetStream not enabled or stream deleted Enable JetStream in config, verify stream exists
Insufficient resources Storage limits exceeded Increase limits or purge old messages
Consumer lag Processing too slow Scale consumers or increase batch size
Message timeout Ack not received in time Increase ack_wait duration

Debugging JetStream:

# Check JetStream is enabled
nats server report jetstream

# View stream details
nats stream info STREAM_NAME

# Check consumer lag
nats consumer info STREAM_NAME CONSUMER_NAME

# Purge stream to clear backlog
nats stream purge STREAM_NAME --force

# View pending/unprocessed counts (shown in consumer info, or as a report)
nats consumer report STREAM_NAME

Performance Issues

Slow Publishing:

# Use async publishing in Go
js.PublishAsync("subject", []byte("data"))

# Batch messages
for i := 0; i < 1000; i++ {
    js.PublishAsync("subject", data)
}
select {
case <-js.PublishAsyncComplete():
    // All published
}

Slow Consumption:

// Increase batch size for pull consumers
msgs, _ := sub.Fetch(100, nats.MaxWait(5*time.Second))

// Use multiple goroutines
for i := 0; i < 10; i++ {
    go func() {
        for msg := range msgChan {
            process(msg)
            msg.Ack()
        }
    }()
}

Memory Issues:

# Monitor memory usage
nats server report jetstream

# Use file storage instead of memory
nats stream add --storage file

# Set appropriate limits
nats stream add --max-bytes 1GB --max-age 7d

Cluster Issues

Issue Symptoms Solution
Split brain Cluster partitioned Check network connectivity between nodes
Slow replication JetStream lag across cluster Check network bandwidth, reduce replication factor
Node not joining Server isolated Verify route configuration and firewall rules

Debugging Cluster:

# Check cluster status
nats server list

# View routes
curl http://localhost:8222/routez

# Check JetStream cluster
nats server report jetstream

# Test inter-node connectivity
nats-server --routes nats://other-server:6222 -DV

Message Acknowledgement Issues

// Explicit ack
msg.Ack()

// Negative ack (redeliver)
msg.Nak()

// Terminate (don't redeliver)
msg.Term()

// In progress (extend ack timeout)
msg.InProgress()

Monitoring and Debugging Tips

# Enable debug logging on server
nats-server -DV

# Monitor all subjects
nats sub ">"

# Trace message flow
nats pub --trace subject "data"

# Benchmark performance (run pub and sub benchmarks separately)
nats bench pub subject --clients 10 --msgs 100000
nats bench sub subject --clients 10 --msgs 100000

# Export metrics for Prometheus
# The server's monitoring port (http_port: 8222) serves JSON endpoints
# (/varz, /jsz, etc.), NOT a Prometheus /metrics endpoint.
# Run the separate prometheus-nats-exporter to scrape those endpoints:
#   https://github.com/nats-io/prometheus-nats-exporter
# prometheus-nats-exporter -varz -jsz=all http://localhost:8222

Best Practices

Subject Design:

  • Use hierarchical subjects: app.service.event.action
  • Keep subjects concise but descriptive
  • Avoid too many tokens (performance impact)
  • Use wildcards judiciously in subscriptions

Error Handling:

// Always handle connection errors
nc, err := nats.Connect(url,
    nats.ErrorHandler(func(nc *nats.Conn, sub *nats.Subscription, err error) {
        log.Printf("Error: %v", err)
    }),
)

// Use context for timeouts
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
msg, err := nc.RequestWithContext(ctx, "subject", data)

Resource Management:

// Always close connections
defer nc.Close()

// Use Drain for graceful shutdown
defer nc.Drain()

// Unsubscribe when done
defer sub.Unsubscribe()

JetStream Tuning:

# Set appropriate ack wait for your workload (the flag is --wait)
nats consumer add --wait 60s

# Limit max deliveries to prevent poison messages
nats consumer add --max-deliver 5

# Use pull consumers for better flow control
nats consumer add --pull