skip to content

You run a reactive change-stream listener as a long-lived Flux in production. What failure modes and design concerns must you address?

level: principalimportance: should knowfreq 25%

answer

  1. Infinite hot Flux, subscribe once at startup
  2. retryWhen backoff + resumeAfter token = gap-free
  3. At-least-once => idempotent handlers
  4. Oplog window > max lag, else resync
  5. Multi-instance: leader-elect / partition / dedupe

basics

~20 s

Handle stream errors with resubscription, persist and resume from tokens for gap-free at-least-once delivery, make handlers idempotent, keep the oplog window larger than max consumer lag, bound backpressure, and coordinate consumers so events are processed once across instances.

solid answer

~50 s

A long-lived change-stream Flux is an infinite, hot pipeline, so treat it like an always-on consumer, not a request. Key concerns: (1) Resilience: the Flux errors on connection loss or primary failover, so wrap in retryWhen with backoff and resubscribe. (2) Delivery guarantee: persist the resume token after processing each event and restart with resumeAfter/startAfter for gap-free delivery, which yields at-least-once, so handlers must be idempotent. (3) Oplog window: if a consumer lags longer than the oplog retention, its token expires and you must resync, so size the oplog for worst-case lag. (4) Backpressure: a hot source cannot be slowed by the DB, so bound flatMap concurrency and choose an overflow policy. (5) Multi-instance: two app replicas both open the stream and each sees every event, so you need leader election, sharded routing, or idempotent+deduplicated writes to avoid double side effects. (6) Observability: track lag, resubscribes, and last token.

code

java · 27 lines
java
@Component
class ChangeStreamRunner {
    private final ReactiveMongoTemplate template;
    private final TokenStore tokens;
    private Disposable sub;
    ChangeStreamRunner(ReactiveMongoTemplate t, TokenStore s){this.template=t;this.tokens=s;}

    @EventListener(ApplicationReadyEvent.class)
    void start() {
        ChangeStreamOptions.ChangeStreamOptionsBuilder o = ChangeStreamOptions.builder()
                .fullDocumentLookup(FullDocument.UPDATE_LOOKUP);
        tokens.last().ifPresent(o::startAfter);   // gap-free resume

        sub = template.changeStream("orders", o.build(), Order.class)
                .concatMap(this::handleIdempotently)   // preserve order, bounded
                .doOnNext(tokens::save)                 // commit token AFTER handling
                .retryWhen(Retry.backoff(Long.MAX_VALUE, Duration.ofSeconds(2)).jitter(0.5))
                .subscribe();
    }

    private Mono<BsonValue> handleIdempotently(ChangeStreamEvent<Order> e) {
        // upsert / dedupe so redelivery is safe; return the resume token
        return applySideEffect(e).thenReturn(e.getResumeToken());
    }

    @PreDestroy void stop() { if (sub != null) sub.dispose(); }
}

go deeper

for a junior

Recognize it is a long-lived stream that can fail and restart.

for a middle

Add retry/resume with tokens and idempotent handling.

for a senior

Reason about oplog window, backpressure on hot sources, and ordering with bounded concurrency.

for a principal

Design the whole CDC pipeline: delivery semantics, multi-instance coordination, resync strategy, observability, and when to front it with a durable log like Kafka.

Running a MongoDB **change stream** as a persistent reactive `Flux` is effectively building an always-on **stream consumer / CDC pipeline**. Several production concerns go beyond the happy-path API. ## Keeping the stream alive **1. It is an infinite, hot pipeline.** Unlike a request-scoped query, this `Flux` should be subscribed once at startup and kept alive for the process lifetime. You need a place to own the subscription (a `@Component` that subscribes in `@PostConstruct`/`ApplicationRunner` and disposes on shutdown) and a `Disposable` to cancel cleanly. **2. Errors and resubscription.** The stream **terminates with an error** on: - connection loss, - replica-set **primary failover**, - network blips, - or cursor invalidation. There is no auto-restart. Wrap with `retryWhen(Retry.backoff(...))` (bounded jittered backoff, effectively infinite attempts) so it re-establishes. On resubscribe you must **resume from the last token**, not from now, or you drop events during the outage. ## Delivery guarantees and history **3. Delivery guarantee = at-least-once (with tokens).** Persist the **resume token** *after* the side effect for each event, and restart via **`resumeAfter(token)`** / **`startAfter(token)`** (the latter also survives INVALIDATE). This gives **gap-free at-least-once**: after a crash you reprocess from the last committed token, so some events may be **redelivered**. Therefore **handlers must be idempotent** (upsert by natural key, dedupe by event id/`_id` + operationType, or track processed positions). Persisting the token *before* processing would instead risk **at-most-once** (lost events on crash). **4. Oplog retention vs consumer lag.** Change streams read the **oplog**, a capped log with a finite time window. If a consumer is down or lagging **longer than the oplog window**, its saved token **ages out** and resumption fails (`ChangeStreamHistoryLost`). Mitigations: - size the oplog for worst-case downtime, - alert on lag, - and have a **resync/bootstrap** path (full scan then resume with a fresh token) for when history is lost. ## Load, scale and ordering **5. Backpressure on a hot source.** Writers keep producing regardless of your demand, so the DB cannot throttle you. If your handler is slower than the write rate you must bound memory: bounded `flatMap(fn, concurrency)`, `limitRate`, and a deliberate **overflow policy** (`onBackpressureBuffer` bounded, `Drop`, or `Latest`). Unbounded buffering risks OOM; dropping risks data loss — this is a real design decision, not a default. **6. Multiple app instances (the big one).** In a horizontally scaled service, **every replica** that opens the stream receives **every event independently**. If each triggers the side effect you get **duplicate work** (double notifications, double index writes). Options: - (a) **single active consumer** via leader election (e.g., a distributed lock / `ShedLock`-style) so only one instance runs the stream; - (b) **partition** the stream (filter by a shard key range per instance); - (c) make all side effects **idempotent and deduplicated** so duplicates are harmless. Choose based on cost of duplicate side effects. **7. Ordering and causal concerns.** A single change stream preserves total order of events; but if you fan out with concurrent `flatMap`, you can reorder processing. If order per entity matters, key-partition (e.g., `groupBy` on document id) so same-entity events process sequentially while different entities go in parallel. ## Payloads, observability and shutdown **8. Update payload completeness.** Without `fullDocument = UPDATE_LOOKUP`, UPDATE events carry only deltas; DELETE events carry only `_id` (unless pre-images are enabled on newer MongoDB). Design downstream logic accordingly. **9. Observability & ops.** Emit metrics: - events/sec, - processing lag (now minus event clusterTime), - resubscribe count, - last committed token age vs oplog window, - error rate. These tell you before history is lost. **10. Graceful shutdown.** On SIGTERM, stop pulling, finish in-flight events, persist the final token, then dispose the subscription — so restart resumes cleanly. ## When this is the right design Change-stream-driven pipelines are excellent for cache/search-index sync, outbox-style event propagation, and notifications. If you need stronger ordering, multi-topic routing, or replay beyond the oplog window, front it with a durable log (Kafka) fed by the change stream, and treat Mongo as the CDC source.

  • Two replicas of your service both open the same change stream. What breaks and how do you fix it?
    Each replica receives every event independently, so any side effect runs twice (duplicate notifications/writes). Fix with a single active consumer via leader election / distributed lock, or partition the stream per instance, or make every side effect idempotent and deduplicated so duplicates are harmless.
  • Should you persist the resume token before or after processing an event, and why?
    After, so that a crash mid-processing re-delivers the event on restart (at-least-once, gap-free) rather than skipping it. Persisting before processing risks at-most-once with lost events. Because after-commit means possible redelivery, handlers must be idempotent.
  • What causes ChangeStreamHistoryLost and how do you recover?
    The saved resume token has aged out of the finite oplog window because the consumer was down or lagging too long. Recovery requires a full resync/bootstrap of the affected data, then starting a fresh stream from a new token; prevention is sizing the oplog above worst-case lag and alerting on lag.

saying these in an interview costs you the question

  • Assuming the Flux auto-restarts after a primary failover without retry logic
  • Persisting the token before processing and calling it exactly-once
  • Ignoring that every scaled instance receives every event (duplicate side effects)
  • Believing the oplog retains history indefinitely so tokens never expire
  • Treating a hot change stream as self-throttling like a normal query

context