NATS JetStream as a Task Queue
NATS is best known as a fast, fire-and-forget messaging system; JetStream adds persistence, acknowledgements, and redelivery, which make it a credible task queue. This guide configures JetStream for background jobs and compares it with the usual brokers, as part of Message Broker Comparison in Queue Fundamentals & Architecture.
Problem Statement
A platform team already runs a NATS cluster for service-to-service messaging and wants to avoid adding RabbitMQ or Redis just for background jobs. The jobs are image processing and outbound webhooks: tens of thousands per hour, bursty, each taking between 100 ms and 30 seconds, with retries for transient failures and a place to park jobs that keep failing. The question is whether JetStream can provide durable storage, at-least-once delivery with redelivery on timeout, work distribution across competing workers, bounded retries, and a dead-letter destination — and what the configuration looks like.
Prerequisites
- NATS Server 2.10+ with JetStream enabled, ideally a 3-node cluster for replicated streams.
- The
natsCLI for administration and a client library with JetStream support (Gonats.gowith thejetstreampackage, Pythonnats-py, Nodenats). - Subject naming conventions for jobs, e.g.
jobs.<kind>. - Idempotent handlers — JetStream delivers at least once.
Step 1 — Create a Stream with WorkQueue Retention
A stream stores messages published to its subjects. With WorkQueue retention, a message is removed as soon as a consumer acknowledges it, and each message can be consumed by only one consumer — the semantics of a task queue rather than a log.
nats stream add JOBS \
--subjects "jobs.>" \
--retention work \ # delete on ack; one consumer per subject filter
--storage file \ # durable on disk
--replicas 3 \ # survive a node loss
--max-age 7d \ # safety net: drop anything unprocessed for a week
--max-msgs-per-subject -1 \
--discard old \
--dupe-window 2m # dedupe publishes with the same Nats-Msg-Id
The duplicate window deduplicates publishes: a producer that retries with the same Nats-Msg-Id header within two minutes does not create a second job. It is JetStream's equivalent of SQS FIFO's deduplication id.
Step 2 — Create a Durable Pull Consumer per Job Kind
Pull consumers let workers request messages when they have capacity, which is the natural fit for competing workers. Each job kind gets its own durable consumer with a subject filter, so kinds can have different timeouts and retry limits.
nats consumer add JOBS images \
--filter "jobs.images" \
--pull \
--ack explicit \ # every message must be acked, nacked, or termed
--wait 60s \ # AckWait: redeliver if not acked within 60 s
--max-deliver 6 \ # 1 attempt + 5 redeliveries
--backoff "10s,30s,2m,10m,30m" \ # per-redelivery delays (overrides AckWait for redeliveries)
--max-pending 2000 \ # cap on unacked messages across all workers
--deliver all \
--replay instant
--backoff gives exponential redelivery spacing without any worker-side scheduling. --max-pending is a fleet-wide flow-control limit: once 2,000 messages are in flight unacknowledged, fetches return nothing until some are acked — a built-in guard against a slow downstream being flooded.
Step 3 — Write a Worker with Explicit Acks and Delayed Naks
Workers fetch in batches and must end each message with one of four outcomes: ack (done), nak (retry, optionally after a delay), in_progress (extend the ack deadline), or term (never redeliver).
// worker.go — nats.go jetstream API
js, _ := jetstream.New(nc)
cons, _ := js.Consumer(ctx, "JOBS", "images")
for {
batch, err := cons.Fetch(10, jetstream.FetchMaxWait(5*time.Second))
if err != nil {
continue
}
for msg := range batch.Messages() {
meta, _ := msg.Metadata()
err := processImage(ctx, msg.Data(), func() { msg.InProgress() }) // heartbeat on long work
switch {
case err == nil:
msg.Ack()
case errors.Is(err, ErrCorruptImage):
msg.Term() // permanent: stop redelivering
deadLetter(js, msg, meta, err)
case errors.Is(err, ErrRateLimited):
msg.NakWithDelay(30 * time.Second) // back off without burning a slot
default:
msg.Nak() // transient: consumer backoff applies
}
}
}
InProgress() resets the ack timer, so a job that legitimately takes 45 seconds under a 60-second AckWait can extend itself while it works — the JetStream form of a visibility-timeout heartbeat, as described in configuring visibility timeouts for long-running workers.
Step 4 — Build a Dead-Letter Destination
JetStream has no built-in dead-letter queue. When a message reaches MaxDeliver, the server publishes an advisory on $JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.<stream>.<consumer> containing the stream sequence. A small service subscribes to advisories, fetches the original message by sequence, and republishes it to a dead-letter stream.
// dlq.go — move max-delivered messages to a DLQ stream
js.CreateStream(ctx, jetstream.StreamConfig{
Name: "JOBS_DLQ", Subjects: []string{"dlq.>"}, Retention: jetstream.LimitsPolicy,
MaxAge: 14 * 24 * time.Hour, Replicas: 3,
})
nc.Subscribe("$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.JOBS.*", func(m *nats.Msg) {
var adv struct{ Stream string `json:"stream"`; Consumer string `json:"consumer"`; StreamSeq uint64 `json:"stream_seq"` }
json.Unmarshal(m.Data, &adv)
stream, _ := js.Stream(ctx, adv.Stream)
orig, err := stream.GetMsg(ctx, adv.StreamSeq)
if err != nil {
return // already removed (e.g., max-age)
}
hdr := nats.Header{"X-Original-Subject": []string{orig.Subject},
"X-Consumer": []string{adv.Consumer}}
js.PublishMsg(ctx, &nats.Msg{Subject: "dlq." + orig.Subject, Data: orig.Data, Header: hdr})
stream.DeleteMsg(ctx, adv.StreamSeq) // remove from the work queue
})
For Term'd messages, the worker publishes to the DLQ itself (the deadLetter call in Step 3), because terminated messages do not generate a max-deliveries advisory. Alert on the DLQ stream's message count as in alerting on dead-letter queue growth.
Step 5 — Compare with the Usual Task-Queue Brokers
| Capability | NATS JetStream | RabbitMQ | Redis (Sidekiq/BullMQ) | SQS |
|---|---|---|---|---|
| Durable, replicated storage | Yes (R3 file streams) | Quorum queues | Depends on persistence config | Managed |
| Redelivery with backoff | Consumer BackOff |
Plugins or delayed retry queues | Framework-level | Visibility timeout |
| Dead-letter queue | Build via advisories | Native DLX | Framework failed set | Native redrive |
| Publish deduplication | Nats-Msg-Id window |
No | Framework job ids | FIFO dedup id |
| Priorities | Separate subjects/consumers | Native per-message | Framework-level | Separate queues |
| Job framework ecosystem | Thin | Celery etc. | Rich | Moderate |
The last row is the practical difference: JetStream gives solid primitives but little of the job-framework tooling (UIs, retry dashboards, scheduling helpers) that Sidekiq, BullMQ, or Celery provide. Teams that already run NATS often find that trade worth it; teams choosing from scratch for background jobs alone usually get further faster with a job framework. The broader decision is covered in Message Broker Comparison.
Step 6 — Monitor Consumers and Scale Workers
The key signals are pending messages (not yet delivered), unacknowledged messages (in flight), and redelivery counts per consumer.
nats consumer info JOBS images --json | jq '{pending: .num_pending, ack_pending: .num_ack_pending, redelivered: .num_redelivered, waiting: .num_waiting}'
# From the NATS Prometheus exporter (surveyor)
nats_consumer_num_pending{stream_name="JOBS", consumer_name="images"}
rate(nats_consumer_delivered_consumer_seq{stream_name="JOBS", consumer_name="images"}[5m])
Scale workers on num_pending (KEDA has a NATS JetStream scaler), and alert when num_ack_pending sits at max-pending — that means workers are saturated or stuck, and fetches are being refused.
Verification
# Publish a job, kill the worker mid-process, confirm redelivery after AckWait
nats pub jobs.images '{"image_id":"i-1"}' -H "Nats-Msg-Id:img-i-1"
nats pub jobs.images '{"image_id":"i-1"}' -H "Nats-Msg-Id:img-i-1" # deduplicated
nats stream info JOBS --json | jq '.state.messages' # expect 1
# Poison message: confirm it reaches the DLQ after 6 deliveries
nats pub jobs.images 'not-json'
sleep 3000; nats stream info JOBS_DLQ --json | jq '.state.messages'
Gotchas & Edge Cases
Overlapping consumers on WorkQueue streams. A WorkQueue stream refuses a second consumer whose subject filter overlaps an existing one; design one consumer per subject or non-overlapping filters.
AckWait vs BackOff. When BackOff is set, it governs redelivery delays and AckWait applies to the first delivery only; set max-deliver to at least the number of backoff entries plus one.
Max-age drops unprocessed jobs. A stream max-age silently discards messages that were never processed. Treat it as a last-resort safety net and alert well before jobs get that old.
Ordering. Messages on one subject are stored in order, but with several workers and redeliveries, processing order is not preserved — the same caveat as any competing-consumers queue.
FAQ
Is JetStream exactly-once?
It offers publish deduplication and double-ack (AckSync) to confirm acknowledgement, which narrows duplicates, but processing is still at-least-once. Keep handlers idempotent.
How do priorities work without per-message priority?
Split priorities into subjects (jobs.images.high, jobs.images.low) with separate consumers, and have workers fetch from the high-priority consumer first, falling back to the low one only when it returns nothing. Because fetch is pull-based, the worker decides the ratio — a simple loop that tries high, then low, gives strict priority; fetching a fixed proportion from each gives weighted fairness without starvation.
What happens to jobs during a NATS node failure?
With --replicas 3, the stream's leader moves to another node and consumers reconnect; acknowledged messages stay acknowledged and unacknowledged ones are redelivered after AckWait. Workers should use the client's reconnect handling and treat a fetch error as "try again", not as a fatal condition. A single-replica stream on a failed node is unavailable until that node returns.
Push or pull consumers for jobs? Pull. Push consumers deliver as fast as the server can, which suits event fan-out; pull lets each worker take only what it can handle.
Can JetStream schedule a job for later?
Not natively per message in most versions; use NakWithDelay for retries, or keep delayed jobs elsewhere (a database or scheduler) and publish when due.
Related
- Message Broker Comparison — where JetStream fits among brokers.
- Kafka vs RabbitMQ for Task Queues — the log vs queue trade-off JetStream straddles.
- Dead-Letter Queues & Poison Messages — DLQ design principles used in Step 4.
- Choosing a Go Job Queue Library — JetStream among Go options.