nats-io/nats-server · warning · errConsumerClosed
consumer closed
Error message
consumer closed
What it means
`errConsumerClosed` is a sentinel error from the JetStream consumer code: it signals that a replicated consumer's `Consumer` object has been closed (e.g. due to leader change, consumer deletion, or shutdown) while its raft group was still being used. Callers deliberately special-case it (server/jetstream_cluster.go:7567, 7950) to avoid logging spurious errors during normal leader transitions. Exported paths that surface it indicate the consumer ceased to exist mid-operation.
Source
Thrown at server/jetstream_cluster.go:8046
}
o.mu.Unlock()
case removePendingRequest:
o.mu.Lock()
if !o.isLeader() {
if o.prm != nil {
delete(o.prm, string(buf[1:]))
}
}
o.mu.Unlock()
default:
return fmt.Errorf("unknown consumer entry op type: %v", entryOp(buf[0]))
}
}
}
return nil
}
var errConsumerClosed = errors.New("consumer closed")
func (o *consumer) processReplicatedAck(dseq, sseq uint64) error {
o.mu.Lock()
// Update activity.
o.lat = time.Now()
var ackAllSeqs []uint64
if o.retention != LimitsPolicy && (o.cfg.AckPolicy == AckAll || o.cfg.AckPolicy == AckFlowControl) {
// Always use the store state, as o.asflr is skipped ahead already.
// Capture before updating store, which clears the pending below.
state, err := o.store.BorrowState()
if err == nil {
// Only need to collect if the ack covers more than the sequence itself.
if sagap := sseq - state.AckFloor.Stream; sagap > 1 {
// At most the pending entries below the ack, don't over-allocate.
ackAllSeqs = make([]uint64, 0, min(uint64(len(state.Pending)), sagap-1))
for seq := range state.Pending {
if seq < sseq {View on GitHub (pinned to 3a66a489d2)
Solutions
- Confirm the consumer still exists with `nats consumer info <stream> <consumer>`; recreate it if it was deleted
- Retry the operation after the new consumer leader stabilizes — the error during leader change is expected and self-healing
- Stop client ack/subscribe loops gracefully before deleting consumers
- Check cluster stability (flapping leader elections) if the error recurs frequently
Defensive patterns
Strategy: retry
Validate before calling
// Verify the consumer still exists before acking/subscribing
_, err := js.ConsumerInfo(stream, consumer)
if err != nil { recreateOrLookupConsumer() } Try / catch
// Treat 'consumer closed' as transient during leader elections
if err := o.processReplicatedAck(dseq, sseq); err != nil {
if errors.Is(err, errConsumerClosed) { /* resync with new leader, do not log fatal */ }
} Prevention
- Stop consumers gracefully before deleting them
- Stabilize the JetStream cluster to avoid frequent leader elections
- Recreate consumers with durable names so clients can resubscribe after closure
When it happens
Trigger: `o.processReplicatedAck(dseq, sseq)` on an already-closed consumer; applying replicated consumer entries to a consumer whose raft group was torn down; snapshot/entry application racing a consumer delete or leadership loss (server/jetstream_cluster.go:8046).
Common situations: Deleting a pull consumer while clients are still acking messages; frequent leader elections in a JetStream cluster causing consumer reassignment; server shutdown or asset rebalancing racing consumer operations.
Related errors
- JS_CONSUMER_OFFLINE
- consumer write error: %v
- unsupported consumer %q
- max_request_batch must be set if it's JetStream limits are s
- consumer name can not contain '.', '*', '>', '\', '/'
AI-assisted analysis of nats-io/nats-server@3a66a489d2 (2026-09-02).
Data as JSON: /api/errors/70f6487e80587f6e.
Report an issue: GitHub.