NATS JetStream

NATS (https://github.com/nats-io/nats-server 20,832 ) (Apache 2.0, a CNCF project) is a single Go binary for subject-based messaging: publishers send to subjects such as booknest.orders.order_paid, subscribers listen with wildcards (* for one token, > for the rest), and core NATS delivers at most once. JetStream, added in 2.2 (March 2021), puts persistence in the same server. A stream captures subjects into files or memory, replicated with Raft; consumers are server-side cursors on it, push or pull, with acknowledgments and redelivery. Retention is limits, interest (until every consumer has acknowledged) or workqueue (until one has), and NATS's key-value and object stores are built on streams. Version 2.15.0 (17 Sep 2026) is current, and this starts it with JetStream on: docker run -d --name l1-nats -p 31422:4222 nats:2.15.0 -js -sd /data.

The Python client (pip 21,050 install nats-py, 2.16.0) creates the stream, publishes each event to a subject named after its type with a Nats-Msg-Id header, and resends the first 500 as a retrying producer would. Two packers share a durable pull consumer that sees only order_paid.

nats_orders.py: deduplicated publishing and a pull consumer on one subjectPython
"""nats_orders.py: 20,000 BookNest events into JetStream, a retry, two pull consumers."""
import asyncio, json
import nats
async def main():
    nc = await nats.connect("nats://localhost:31422")
    js = nc.jetstream()
    await js.add_stream(name="ORDERS", subjects=["booknest.orders.>"],
                        max_age=7 * 86400, duplicate_window=120)        # seconds
    with open("data/order_events.jsonl") as f:
        events = [json.loads(line) for line, _ in zip(f, range(20_000))]
    def publish(e):                  # one subject per event type; the id makes retries safe
        return js.publish(f"booknest.orders.{e['type']}", json.dumps(e).encode(),
                          headers={"Nats-Msg-Id": str(e["event_id"])})
    for i in range(0, 20_000, 1_000):
        await asyncio.gather(*map(publish, events[i:i + 1_000]))
    retry = await asyncio.gather(*map(publish, events[:500]))      # a producer retries
    print("retried 500, stored again:", sum(not ack.duplicate for ack in retry))
    sub = await js.pull_subscribe("booknest.orders.order_paid", durable="packers")
    async def packer(done=0):
        while True:
            try:
                batch = await sub.fetch(50, timeout=1)
            except nats.errors.TimeoutError:
                return done
            done += len(batch)
            await asyncio.gather(*(m.ack() for m in batch))
    print("packed:", await asyncio.gather(packer(), packer()))
    print("pending:", (await js.consumer_info("ORDERS", "packers")).num_pending)
    await nc.close()
asyncio.run(main())
Output
retried 500, stored again: 0
packed: [2583, 2582]
pending: 0

The server dropped all 500 retries inside the two-minute duplicate window, the packers shared the 5,165 paid orders, and none is pending. The server used under 20 MiB, which suits edge sites linked to a central cluster by leaf nodes. The limits are scale and ecosystem: one Raft leader takes all of a stream's writes, so you spread load across streams yourself, and nothing matches Kafka Connect 129 . The licence had a scare: in April 2025 Synadia, NATS's main sponsor, announced a move to the Business Source License, then agreed with the CNCF on 1 May 2025 that NATS stays there under Apache 2.0.