Using NATS from asyncio with nats-py¶
NATS is a small, fast messaging server, and nats-py is its asyncio client. Core NATS gives you publish/subscribe and request/reply with no persistence: a message goes to whoever is subscribed at that moment, or nowhere. JetStream, built into the same server, adds streams that store messages and consumers that acknowledge them. Measured with nats-py 2.16 against nats-server 2.15 on the same host: a request/reply round trip took 155 µs; core publish/subscribe moved 224,204 100-byte messages per second; and a subscriber whose callback took 1 ms per message received 1,006 of 20,000 — the other 18,994 were dropped as a slow consumer, reported only through the error callback. JetStream published 8,581 msg/s awaiting each acknowledgement and 54,669 with publish_async, and a pull consumer acknowledged 80,157 msg/s. This guide uses each mode for what it is good at.
Prerequisites¶
- Python 3.11+,
pip install nats-py; anats-serverstarted with-jsfor JetStream. - Broker trade-offs, from Message Brokers & Event Streams.
- Callback concurrency, from Task Scheduling & Lifecycle.
1. Connect with reconnect and error callbacks¶
One connection per process carries every subscription and publish. Configure what happens when the server goes away, and always register an error callback — it is where NATS reports problems that do not raise:
import nats
async def connect() -> nats.NATS:
async def on_error(exc: Exception) -> None:
log.error("nats error: %s", exc) # slow consumers are reported here
async def on_disconnect() -> None:
log.warning("nats disconnected")
return await nats.connect(
servers=["nats://nats-1:4222", "nats://nats-2:4222"],
name="billing-api",
max_reconnect_attempts=-1, # keep trying forever
reconnect_time_wait=2,
error_cb=on_error,
disconnected_cb=on_disconnect,
)
While disconnected, nats-py buffers publishes up to a limit and flushes them after reconnecting; subscriptions are re-established automatically. Close the connection with await nc.drain() at shutdown, which stops subscriptions, processes what is already received, flushes publishes, and then closes — the graceful form of close().
Verify: restart the NATS server under light traffic; the client logs a disconnect and a reconnect, and subscriptions resume without code changes.
2. Use request/reply for RPC-style calls¶
Request/reply sends a message with a unique reply subject and waits for the first answer. Services subscribe with a queue group so requests are spread across instances:
async def ping_handler(msg) -> None:
await msg.respond(b"pong:" + msg.data)
await nc.subscribe("svc.ping", cb=ping_handler, queue="ping-workers") # load-balanced group
reply = await nc.request("svc.ping", b"hello", timeout=1.0)
print(reply.data) # b"pong:hello"
Measured: 155 µs per sequential round trip over localhost. A request to a subject nobody subscribes to fails at once with NoRespondersError — tested — instead of waiting for the timeout, which makes "service not running" distinguishable from "service slow". Queue groups give load balancing without any configuration: each request goes to one member of the group, and instances can join or leave at any time.
Verify: with two service instances in the same queue group, requests are split between them, and stopping both produces NoRespondersError rather than timeouts.
3. Keep core subscribers fast, or lose messages¶
A core subscription has a bounded buffer of pending messages. When the callback cannot keep up and the buffer fills, NATS drops messages for that subscription and reports a slow consumer through the error callback — nothing raises in your handler:
async def slow_handler(msg) -> None:
await asyncio.sleep(0.001) # 1 ms of work per message
await nc.subscribe("events", cb=slow_handler, pending_msgs_limit=1000)
# measured: 20,000 published, 1,006 received, 18,994 dropped as slow consumer
Callbacks for one subscription run one at a time, so a 1 ms callback caps that subscription at about 1,000 messages per second regardless of how many cores you have. Either keep callbacks short and hand work to a bounded worker pool, or — when every message matters — use JetStream, where unacknowledged messages are redelivered instead of dropped:
work: asyncio.Queue = asyncio.Queue(maxsize=10_000)
async def enqueue(msg) -> None:
work.put_nowait(msg) # raises QueueFull instead of silently losing messages
Verify: under peak publish rate, the error callback reports no slow-consumer events.
4. Use JetStream when messages must not be lost¶
JetStream stores messages in a stream and tracks delivery per consumer. Publishing returns an acknowledgement from the server once the message is stored:
js = nc.jetstream()
await js.add_stream(name="ORDERS", subjects=["orders.>"])
ack = await js.publish("orders.created", payload) # stored: ack.seq is its position
# Higher throughput: send many, then wait for all acknowledgements
futures = [await js.publish_async("orders.created", p) for p in payloads]
await asyncio.gather(*futures)
Measured: 8,581 msg/s awaiting each publish acknowledgement, and 54,669 msg/s with publish_async — the same "don't wait per message" rule as for Kafka producers in publishing to Kafka with aiokafka producers. Configure the stream's retention (limits, interest or work-queue), storage (file or memory) and replicas to match how long and how safely messages must be kept.
Verify: stop all consumers, publish, restart them; every message published while they were down is delivered.
5. Consume with durable pull consumers¶
A durable consumer remembers its position across restarts. Pull consumers let the application ask for batches when it has capacity, which is natural backpressure for asyncio workers:
sub = await js.pull_subscribe("orders.>", durable="billing")
async def run() -> None:
while True:
try:
batch = await sub.fetch(100, timeout=5)
except nats.errors.TimeoutError:
continue # nothing new
for msg in batch:
try:
await process(msg.data)
await msg.ack()
except PermanentError:
await msg.term() # never redeliver
except Exception:
await msg.nak(delay=5) # redeliver after 5 s
Measured: fetching 500 at a time and acknowledging each, 80,157 msg/s; after consuming 5,000, the consumer reported an ack floor of 5,000 and the remaining messages pending. msg.term() stops redelivery of a poison message and msg.nak(delay=...) schedules a retry; set max_deliver on the consumer to bound attempts, as discussed in handling poison messages in async consumers. Several processes sharing the same durable name share the work.
Verify: kill a consumer mid-batch; unacknowledged messages are redelivered to the next fetch after the ack wait.
Verification¶
NATS is used well when:
- One connection per process, with reconnects, an error callback and
drain()at shutdown. - Request/reply services use queue groups, and clients distinguish no responders from timeouts.
- Core subscribers never trigger slow-consumer events at peak.
- Messages that must not be lost go through JetStream with durable consumers and bounded redelivery.
Diagnostic Hook: count slow-consumer errors from the error callback and compare published with received counts per subject. Any slow-consumer event on a core subscription is silent message loss; JetStream consumer num_pending and num_redelivered show backlog and retry pressure.
Pitfalls & edge cases¶
- No error callback. Slow-consumer drops are invisible: 18,994 of 20,000 lost in testing.
- Slow work in core callbacks. Callbacks run one at a time per subscription.
- Core NATS for durable work. Messages published with no subscriber are gone.
- Awaiting each JetStream publish. 8,581 msg/s against 54,669 with
publish_async.
Frequently Asked Questions¶
How do I use NATS with asyncio?
Connect with await nats.connect(...), subscribe with a coroutine callback, publish with await nc.publish(subject, data), and use await nc.request(subject, data, timeout=...) for request/reply. Use nc.jetstream() for persistent streams.
Why is my NATS subscriber missing messages?
It is probably a slow consumer: when the callback falls behind and the pending buffer fills, NATS drops messages and reports it only through the error callback. In testing a 1 ms callback received 1,006 of 20,000.
What is the difference between core NATS and JetStream?
Core NATS delivers only to current subscribers and stores nothing. JetStream stores messages in streams and tracks acknowledgements per consumer, so messages survive consumer restarts.
How fast is NATS request/reply from Python?
About 155 µs per round trip on localhost in testing, and a request with no subscribers fails immediately with NoRespondersError.
Related¶
- Message Brokers & Event Streams — up to the topic overview.
- Reading Redis Streams with consumer groups — another lightweight option for durable consumption.
- Network I/O & Protocol Handling — the section overview.