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.
graph TB
subgraph "NATS Architecture"
subgraph "NATS Server Cluster"
N1[NATS Server 1]
N2[NATS Server 2]
N3[NATS Server 3]
N1 <-->|Cluster| N2
N2 <-->|Cluster| N3
N3 <-->|Cluster| N1
end
subgraph "Publishers"
P1[Service A]
P2[Service B]
end
subgraph "Subscribers"
S1[Worker 1]
S2[Worker 2]
S3[Dashboard]
end
subgraph "JetStream"
JS[(Stream Storage)]
end
P1 -->|Publish| N1
P2 -->|Publish| N2
N1 -->|Subscribe| S1
N2 -->|Subscribe| S2
N1 -->|Subscribe| S3
N1 <-->|Persist| JS
N2 <-->|Persist| JS
N3 <-->|Persist| JS
end
style N1 fill:#e1f5fe
style N2 fill:#e1f5fe
style N3 fill:#e1f5fe
style JS fill:#fff3e0
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.*.londonmatchessensors.temperature.london)>- Matches one or more tokens at the end (e.g.,orders.>matchesorders.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)
flowchart LR
P[Publisher] -->|Publish to 'orders.created'| N[NATS Server]
N -->|Deliver| S1[Subscriber 1]
N -->|Deliver| S2[Subscriber 2]
N -->|Deliver| S3[Subscriber 3]
style N fill:#e1f5fe
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
flowchart LR
P[Publisher] -->|orders.process| N[NATS Server]
subgraph "Queue Group 'workers'"
W1[Worker 1]
W2[Worker 2]
W3[Worker 3]
end
N -->|Msg 1| W1
N -->|Msg 2| W3
N -->|Msg 3| W2
N -->|Msg 4| W1
style N fill:#e1f5fe
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
graph TB
subgraph "JetStream Architecture"
P[Publisher] -->|Publish| S[Stream]
S[(Stream<br/>orders.created<br/>Subjects: orders.*)] -->|Pull/Push| C1[Consumer 1<br/>Durable]
S -->|Pull/Push| C2[Consumer 2<br/>Ephemeral]
S -->|Pull/Push| C3[Consumer 3<br/>Queue Group]
C1 -->|Ack| S
C2 -->|Ack| S
C3 -->|Ack| S
S -.->|Retention Policy| R[Retention<br/>Limits/Age/Interest]
end
style S fill:#fff3e0
style C1 fill:#e8f5e9
style C2 fill:#e8f5e9
style C3 fill:#e8f5e9
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)
stateDiagram-v2
[*] --> Pending: Message arrives
Pending --> Delivered: Consumer fetches
Delivered --> Acknowledged: Consumer acks
Delivered --> Redelivered: Ack timeout
Redelivered --> Delivered: Retry
Redelivered --> MaxRedeliveries: Too many retries
Acknowledged --> [*]
MaxRedeliveries --> [*]: Move to dead letter
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
graph TB
subgraph "Cluster 1 - EU"
N1[NATS EU-1]
N2[NATS EU-2]
N3[NATS EU-3]
N1 <-->|Route| N2
N2 <-->|Route| N3
N3 <-->|Route| N1
end
subgraph "Cluster 2 - US"
N4[NATS US-1]
N5[NATS US-2]
N4 <-->|Route| N5
end
N1 <-.->|Gateway| N4
L1[Leaf Node] -.->|Leaf| N1
style N1 fill:#e1f5fe
style N2 fill:#e1f5fe
style N3 fill:#e1f5fe
style N4 fill:#c8e6c9
style N5 fill:#c8e6c9
style L1 fill:#fff9c4
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 <...>(beyondconnection), andnats server report ...issue requests on the$SYSsystem 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 asystem_accountin the server config.nats server check connection,nats rtt, andnats account infowork 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