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

  1. Confirm the consumer still exists with `nats consumer info <stream> <consumer>`; recreate it if it was deleted
  2. Retry the operation after the new consumer leader stabilizes — the error during leader change is expected and self-healing
  3. Stop client ack/subscribe loops gracefully before deleting consumers
  4. 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

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


AI-assisted analysis of nats-io/nats-server@3a66a489d2 (2026-09-02). Data as JSON: /api/errors/70f6487e80587f6e. Report an issue: GitHub.