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.

Advertisement

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:

SubjectDirectionPurpose
telemetry.<site>.<device>.<metric>upPeriodic readings, captured by a stream
event.<site>.<device>.<kind>upAlarms and state changes, longer retention
cmd.<site>.<device>.<action>downCommands, one consumer per device
rpc.<device>.<method>down, reply upLive 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.

Advertisement

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.

Site leaf nodes buffer locally; the hub cluster sources their streamsSensorsNATS client, TLSLegacy devicesMQTT 3.1.1Actuatorspull cmd consumerLeaf node (site)JetStream domain site17stream TELEMETRY_S17stream COMMANDS_S17publishmqttcommandsHub cluster (R3)JetStream domain hubTELEMETRY: sources onlypublishes cmd.s17... to the siteKV device_stateleafnode linkcommandsTSDB writerpull, batch ackRules enginepull, filteredWAN outageleaf keeps writingOn reconnect the hub source resumes from its last sequenceSubjects: telemetry.site.device.metric and cmd.site.device.action
Devices publish to a leaf node at their site. Site streams buffer during WAN loss; the hub cluster sources them into one TELEMETRY stream and sends commands back down.

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))
StreamRetentionKey limitsWhy
TELEMETRYlimitsmax age 7 days, max messages per subjectA replayable window; one chatty device cannot evict everyone else
EVENTSlimitsmax age 90 days, R3 at the hubAlarms feed audits and incident reviews
COMMANDSwork queuemax age 1 day, per-message TTLA message is deleted once its device acknowledges it
device_state (KV)key-value buckethistory 1 to 5Last 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

SymptomCauseFix
Duplicate readings after flaky radioMessage id regenerated per attemptId from device id plus persisted sequence
Stale command executed hours laterNo TTL, no expiry check on the deviceNats-TTL on commands, plus a not-after timestamp the device checks
Redelivery storm, writer never catches upack_wait shorter than the batch writeRaise ack_wait, shrink batches, watch redelivery counts
One device fills the streamOnly global limits setMax messages per subject, rate limits in device credentials
Hub lag grows after an outageAll sites catching up at onceSize hub ingest for recovery, not steady state
Device reads another device's commandsWildcard permissionsPer-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

  1. Write the subject hierarchy down, with site before device, and derive permissions from it.
  2. Deploy one leaf node at a pilot site with its own JetStream domain, and pull the uplink to confirm devices keep publishing.
  3. Create separate telemetry, events and commands streams, with per-subject limits on telemetry.
  4. Persist a per-device sequence and send it as Nats-Msg-Id on every publish.
  5. Enable per-message TTL on the commands stream, make command handlers idempotent and report results on an event subject.
  6. Load-test recovery: hold every site offline for a day, reconnect them together and measure hub catch-up time.
Key takeaway: NATS JetStream suits a device fleet because one binary covers fast pub/sub, durable replicated streams, an MQTT listener and leaf nodes at the edge. Design the subject tree first, with site before device, and hang streams, consumer filters and permissions off it. Put a leaf node with its own JetStream domain at each site so outages buffer locally, deduplicate publishes with a stable message id, expire commands with per-message TTL, acknowledge only after durable writes and size the hub for recovery rather than steady state.