An IoT backend has two jobs that pull in opposite directions. Telemetry flows up from thousands of devices, constantly, and losing a reading now and then is tolerable. Commands flow down to one specific device, rarely, and must arrive exactly once or not at all, and never hours late. Sites lose their uplink, devices reboot mid-publish, and the cloud side has to keep up with all of it.
NATS fits this shape unusually well because one small server binary covers the plain publish-subscribe layer, durable streams through JetStream, a built-in MQTT listener and leaf nodes that run at the edge. This article builds a fleet backbone on it, from subject naming to sizing. For the MQTT protocol itself see the MQTT article; for choosing between brokers in general, the message broker guide.
Two layers: core NATS and JetStream
Core NATS is a subject-based message router with at-most-once delivery. A publisher sends to a subject such as telemetry.s17.pump042.env; every subscriber whose interest matches receives it; if nobody is listening the message is simply gone. There is no disk on the hot path, which is why it is fast and why it is the right transport for request-reply and live dashboards.
JetStream is the persistence layer built into the same server. A stream captures every message published on its subjects into a replicated log, with limits on age, size and count. Consumers are server-side cursors over a stream: durable, acknowledged, redelivered after a timeout, and filterable by subject. Publishing into a stream returns an acknowledgement with the stored sequence number, so a device knows its reading is on disk before it discards it locally.
Most IoT designs use both. Streams hold anything you cannot afford to lose or must replay, and core subjects carry ephemeral traffic such as a live video thumbnail or a request to read a sensor right now.
Subject design is your schema
Subjects are dot-separated tokens, * matches exactly one token and > matches one or more trailing tokens. Every routing decision, stream membership, consumer filter and permission is written against them, so fix the hierarchy before writing code. A layout that works for most fleets:
| Subject | Direction | Purpose |
|---|---|---|
telemetry.<site>.<device>.<metric> | up | Periodic readings, captured by a stream |
event.<site>.<device>.<kind> | up | Alarms and state changes, longer retention |
cmd.<site>.<device>.<action> | down | Commands, one consumer per device |
rpc.<device>.<method> | down, reply up | Live request-reply, core NATS only |
Put the site before the device, because sites are the unit of network failure and of edge deployment: telemetry.s17.> is everything one leaf node must hold. Keep tokens free of dots and whitespace, use stable device identifiers rather than hostnames, and never encode values that change, such as firmware version, into a subject; put them in headers or the payload.
Permissions follow the same tree. With decentralised JWT authentication each device gets credentials that allow publishing only on telemetry.s17.pump042.> and event.s17.pump042.>, and consuming only its own command subjects. A compromised device can then lie about itself but cannot impersonate a neighbour or read another device's commands.
Architecture: a leaf node per site
A leaf node is a NATS server that connects outward to a hub and extends its subject space. Devices talk to the leaf on the local network, so they keep working when the uplink fails. Give each leaf its own JetStream domain and the site's streams live on local disk: during a WAN outage readings keep landing in the site stream, and when the link returns the hub, which sources from every site stream, resumes from the last sequence it copied. No device-side buffering logic is needed beyond surviving a leaf restart.
The leaf configuration is short. Treat the block below as a sketch to check against the documentation for your server version, especially the domain and remote options:
# leaf node at one site (sketch: check key names against your server version)
server_name: site17-leaf
jetstream {
store_dir: "/var/lib/nats"
domain: "site17" # separate JetStream domain, survives WAN loss
max_file_store: 200GB
}
leafnodes {
remotes = [
{ url: "tls://hub.example.net:7422", credentials: "/etc/nats/site17.creds" }
]
}
mqtt {
port: 8883 # built-in MQTT listener for legacy devices
tls { cert_file: "/etc/nats/tls.crt", key_file: "/etc/nats/tls.key" }
}On the hub, the fleet-wide stream is defined with sources pointing at each site stream through that site's domain API prefix, so the copy is pull-based, resumable and visible as source lag in nats stream info. Give the hub stream sources only and no subjects of its own; otherwise leaf traffic crossing the leafnode link is captured twice, once directly and once through the source.
Commands take the opposite path with no hub stream at all: the issuer publishes cmd.s17.pump042.setpoint with a JetStream publish, the subject crosses the leafnode link and the site's work-queue stream stores it and acknowledges. If the site is offline the publish fails at once, so the issuer knows the command was not queued instead of discovering hours later that it ran late.
Streams and consumers for each traffic class
Each traffic class gets its own stream because they want different retention. The setup code uses the nats-py client; the same options exist in every official client and in the nats CLI.
import nats
from nats.js.api import (StreamConfig, ConsumerConfig, RetentionPolicy,
StorageType, DiscardPolicy, AckPolicy)
async def setup(js):
# Telemetry: a time-bounded log, newest wins when full
await js.add_stream(StreamConfig(
name="TELEMETRY_S17", subjects=["telemetry.s17.>"],
retention=RetentionPolicy.LIMITS, storage=StorageType.FILE,
max_age=7 * 24 * 3600, # seconds in nats-py
max_msgs_per_subject=20_000, # caps one chatty device
discard=DiscardPolicy.OLD, duplicate_window=120, num_replicas=1))
# Commands: each message consumed once by its device
await js.add_stream(StreamConfig(
name="COMMANDS_S17", subjects=["cmd.s17.>"],
retention=RetentionPolicy.WORK_QUEUE, storage=StorageType.FILE,
max_age=24 * 3600, duplicate_window=600, num_replicas=1))| Stream | Retention | Key limits | Why |
|---|---|---|---|
| TELEMETRY | limits | max age 7 days, max messages per subject | A replayable window; one chatty device cannot evict everyone else |
| EVENTS | limits | max age 90 days, R3 at the hub | Alarms feed audits and incident reviews |
| COMMANDS | work queue | max age 1 day, per-message TTL | A message is deleted once its device acknowledges it |
| device_state (KV) | key-value bucket | history 1 to 5 | Last known value per device for dashboards |
A work-queue stream requires that consumer filters do not overlap, which per-device filters such as cmd.s17.pump042.> satisfy naturally. For the last-known value, the key-value store is a stream that keeps only the last few messages per subject; a dashboard reads one key instead of scanning telemetry.
Publishing from a device without duplicates
Devices retry. A publish whose acknowledgement was lost to a radio drop will be sent again, and without care the stream stores it twice. JetStream deduplicates on the Nats-Msg-Id header within the stream's duplicate window, so derive the id from something stable, the device id plus a sequence counter persisted across reboots, never a random value generated per attempt.
import json, nats
async def device_loop(dev_id, read_sensor, apply_command):
nc = await nats.connect("tls://leaf.site17.local:4222",
user_credentials=f"/etc/device/{dev_id}.creds",
max_reconnect_attempts=-1)
js = nc.jetstream()
seq = load_persisted_seq() # survives reboot
cmds = await js.pull_subscribe(f"cmd.s17.{dev_id}.>", durable=f"dev-{dev_id}",
stream="COMMANDS_S17")
while True:
seq += 1
reading = read_sensor()
await js.publish(f"telemetry.s17.{dev_id}.env", json.dumps(reading).encode(),
headers={"Nats-Msg-Id": f"{dev_id}-{seq}"}) # dedup on retry
save_seq(seq)
try:
for m in await cmds.fetch(10, timeout=0.5):
apply_command(json.loads(m.data)) # must be idempotent
await m.ack()
except nats.errors.TimeoutError:
pass
await sleep_until_next_tick()The device also drains its own durable command consumer every tick. Pull consumers suit devices well: the device asks for work when it is ready, so a device that sleeps for an hour simply finds its commands waiting, subject to their TTL.
Commands that expire instead of arriving late
The dangerous failure in device control is not loss but lateness: a valve command queued during an outage and delivered six hours later can be worse than none. JetStream has had per-message TTL since server 2.11. A stream must opt in with AllowMsgTTL, and the switch is one-way, so turn it on deliberately; each message then carries a Nats-TTL header and is removed when it expires, while messages without one fall back to the stream's max age.
nats stream edit COMMANDS_S17 --allow-msg-ttl # one-way switch: cannot be turned off
nats pub cmd.s17.pump042.setpoint '{"rpm": 1200, "id": "c-8812"}' \
--header "Nats-TTL:5m" --header "Nats-Msg-Id:c-8812"Give every command an id that the device records after applying it, so a redelivery after a crash between apply and ack is ignored. Report the outcome on event.<site>.<device>.cmd_result so the issuer learns whether the command ran, expired or failed, rather than guessing from silence.
The cloud side: ingest and back-pressure
Consumers on the hub turn the stream into rows. Fetch in batches, write the batch, and only then acknowledge, so a crash replays the batch instead of losing it:
async def tsdb_writer(js, tsdb):
sub = await js.pull_subscribe(
"telemetry.>", durable="tsdb-writer", stream="TELEMETRY",
config=ConsumerConfig(ack_policy=AckPolicy.EXPLICIT, ack_wait=30,
max_ack_pending=5_000, max_deliver=5))
while True:
try:
msgs = await sub.fetch(500, timeout=2)
except nats.errors.TimeoutError:
continue
rows = [decode(m.subject, m.data) for m in msgs]
try:
await tsdb.write_batch(rows) # idempotent upsert on (device, ts)
except TransientError:
for m in msgs:
await m.nak(delay=5) # back off, redeliver later
continue
for m in msgs:
await m.ack()Three settings carry the load. max_ack_pending caps unacknowledged messages and is your back-pressure valve: when the time-series database slows, the consumer stops receiving rather than ballooning memory. ack_wait must exceed your slowest batch write, or messages are redelivered while still in flight. max_deliver bounds poison messages; watch the advisory the server emits when it is exhausted and park those payloads for inspection. Because redelivery happens, the write must be idempotent, typically an upsert keyed on device and timestamp. The back-pressure article covers the general pattern.
Legacy devices through the MQTT listener
Many sensors only speak MQTT. The NATS server can expose an MQTT 3.1.1 listener that requires JetStream to be enabled, because sessions and retained messages are stored in streams. Topic separators are translated, so an MQTT publish to telemetry/s17/meter9/kwh appears to NATS subscribers as telemetry.s17.meter9.kwh and lands in the same stream as native traffic.
Use QoS 1 for telemetry and confirm exactly which QoS levels and MQTT features your server version supports before relying on them; MQTT 5 features should not be assumed. Keep MQTT topic names inside the same hierarchy, and avoid characters that do not survive the mapping, such as dots inside MQTT topic levels.
Worked example: sizing a 10,000-device fleet
Take 10,000 sensors across 40 sites, each publishing a 200-byte reading every 10 seconds. That is 1,000 messages per second fleet-wide and about 200 KB per second of payload. Per day: 200 KB times 86,400 seconds is roughly 17.3 GB. Seven days of retention is about 121 GB per replica, and the hub stream at three replicas needs about 363 GB of disk across the cluster, before per-message overhead of headers, subjects and index, which for small messages can add a large fraction; measure it with a day of real traffic rather than trusting the payload figure.
Each site holds 1/40th: 25 messages per second and about 3 GB per week, so a leaf on modest hardware can buffer a multi-day outage. On reconnect after a 24-hour outage, a site has about 2.16 million messages to forward; at a sustained catch-up rate of a few thousand per second that drains in minutes, but forty sites recovering at once from a regional outage is the case to load-test, because the hub's ingest consumers see a sudden backlog.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Duplicate readings after flaky radio | Message id regenerated per attempt | Id from device id plus persisted sequence |
| Stale command executed hours later | No TTL, no expiry check on the device | Nats-TTL on commands, plus a not-after timestamp the device checks |
| Redelivery storm, writer never catches up | ack_wait shorter than the batch write | Raise ack_wait, shrink batches, watch redelivery counts |
| One device fills the stream | Only global limits set | Max messages per subject, rate limits in device credentials |
| Hub lag grows after an outage | All sites catching up at once | Size hub ingest for recovery, not steady state |
| Device reads another device's commands | Wildcard permissions | Per-device credentials scoped to its own subjects |
Operating it
Watch four numbers: consumer pending and ack-pending counts, redelivery rate, stream bytes against limits and source lag at the hub. The nats stream report and nats consumer report commands show them, and the server exposes the same data on its monitoring endpoint for Prometheus exporters. Run hub streams at three replicas on separate failure domains; site streams usually run at one replica on one box, because the hub holds the durable copy once forwarded.
Newer servers add features worth knowing: atomic batch publish, distributed counters and delayed message scheduling arrived in 2.12, and 2.14 (April 2026) added fast batch publishing and recurring schedules. Check the release notes for your version before designing around them, and keep a fallback for older leaf nodes still in the field. For offline handling on the device itself, the offline queue article covers the client-side pattern.
What to do next
- Write the subject hierarchy down, with site before device, and derive permissions from it.
- Deploy one leaf node at a pilot site with its own JetStream domain, and pull the uplink to confirm devices keep publishing.
- Create separate telemetry, events and commands streams, with per-subject limits on telemetry.
- Persist a per-device sequence and send it as Nats-Msg-Id on every publish.
- Enable per-message TTL on the commands stream, make command handlers idempotent and report results on an event subject.
- Load-test recovery: hold every site offline for a day, reconnect them together and measure hub catch-up time.