skip to content

How do reactive change streams work in Spring Data MongoDB, and how do they differ from @Tailable?

level: seniorimportance: should knowfreq 40%

answer

  1. changeStream() -> Flux<ChangeStreamEvent>
  2. Oplog-backed; replica set required
  3. Any collection, all op types
  4. Resume token = gap-free restart
  5. UPDATE_LOOKUP for full post-image

basics

~20 s

changeStream on ReactiveMongoTemplate opens a Flux of change events (insert, update, delete, replace) on any collection, backed by the replica-set oplog. Unlike @Tailable it needs no capped collection, sees all operation types, and supports resume tokens.

solid answer

~40 s

Spring Data exposes MongoDB change streams reactively via ReactiveMongoTemplate.changeStream(...), returning a Flux<ChangeStreamEvent<T>>. Each event carries the operation type (INSERT/UPDATE/REPLACE/DELETE/INVALIDATE), the full document (or the delta), and a resume token. Change streams are powered by the replica-set/sharded-cluster oplog, so they require a replica set (even single-node) but work on ordinary collections, not just capped ones. You can filter with an aggregation pipeline (for example match on operationType) and request fullDocument lookup for updates. The killer feature over @Tailable is the resume token: persist it and pass it via resumeAfter/startAfter so after a crash you continue exactly where you left off without missing events. @Tailable, by contrast, only tails inserts on a capped collection and cannot resume. Change streams are the right tool for reliable event-driven integration and cache invalidation.

code

java · 20 lines
java
@Service
class UserChangeConsumer {
    private final ReactiveMongoTemplate template;
    private final TokenStore tokens; // persists last resume token
    UserChangeConsumer(ReactiveMongoTemplate t, TokenStore s){this.template=t;this.tokens=s;}

    Flux<ChangeStreamEvent<User>> listen() {
        ChangeStreamOptions.ChangeStreamOptionsBuilder opts =
                ChangeStreamOptions.builder()
                        .fullDocumentLookup(FullDocument.UPDATE_LOOKUP)
                        .filter(Aggregation.newAggregation(
                                Aggregation.match(
                                        Criteria.where("operationType").in("insert","update"))));
        tokens.last().ifPresent(opts::resumeAfter);

        return template.changeStream("users", opts.build(), User.class)
                .doOnNext(e -> tokens.save(e.getResumeToken()))
                .retryWhen(Retry.backoff(Long.MAX_VALUE, Duration.ofSeconds(2)));
    }
}

go deeper

for a junior

Know change streams emit real-time change events on a collection.

for a middle

List the operation types and the replica-set requirement.

for a senior

Design resumable consumers with tokens, UPDATE_LOOKUP, filtering, and retry; contrast with @Tailable.

for a principal

Architect CDC/integration around oplog window sizing, at-least-once semantics, idempotency, and resync strategy.

A **change stream** is MongoDB's built-in change-data-capture feed. It lets a client subscribe to a real-time stream of data changes on a collection, database, or whole deployment, reading from the replica set's **oplog** (operations log). Spring Data exposes it reactively. **API.** `ReactiveMongoTemplate.changeStream(...)` returns a `Flux<ChangeStreamEvent<T>>`. A fluent form: ``` reactiveMongoTemplate.changeStream(User.class) .watchCollection("users") .filter(Aggregation.newAggregation( match(where("operationType").is("insert")))) .listen(); ``` Each **`ChangeStreamEvent<T>`** gives you: - `getOperationType()` — INSERT, UPDATE, REPLACE, DELETE, INVALIDATE, etc. - `getBody()` — the mapped domain object (the full document; for UPDATE you can request `fullDocument = UPDATE_LOOKUP` to get the post-image, otherwise you get only the changed fields / none). - `getResumeToken()` — an opaque position marker. - `getRaw()` — the underlying driver `ChangeStreamDocument`. **Requirements.** - A **replica set** (or sharded cluster). A standalone `mongod` cannot serve change streams; even local dev needs a single-node replica set. (Contrast: `@Tailable` needs a **capped collection** but not a replica set.) - Change streams work on **normal collections** — no capping needed. **Filtering.** You pass an **aggregation pipeline** (`$match`, `$project`, ...) so the server only sends events you care about, e.g. only `update` operations or only documents in a certain state. **Resumability (the big differentiator).** Every event has a **resume token**. Persist the last processed token; on restart pass it with **`resumeAfter(token)`** (resume *after* that event) or **`startAfter(token)`** (also survives an INVALIDATE). This gives **at-least-once** delivery with no gap, which `@Tailable` cannot provide. If the oplog has rolled past your token (token too old), resumption fails and you must do a full resync. **Listeners.** Beyond ad-hoc `Flux`, Spring Data offers a message-listener abstraction. In the reactive world you typically just keep a subscribed `Flux` alive (subscribe at startup, resubscribe on error). In the blocking world there is `MessageListenerContainer` with `ChangeStreamRequest`; the reactive analog is the `changeStream(...).listen()` `Flux` you manage yourself. **`@Tailable` vs change streams.** | | `@Tailable` | Change stream | |---|---|---| | Collection | must be **capped** | any collection | | Deployment | standalone OK | **replica set required** | | Operations seen | **inserts only** | insert/update/replace/delete/invalidate | | Resume after failure | no token, may miss | **resume token**, no gap | | Filtering | query predicate | aggregation pipeline | | Use for | lightweight append-only tail | reliable CDC / integration | **Gotchas.** (1) Without `UPDATE_LOOKUP`, UPDATE events carry only the delta (`updateDescription`), not the whole document — mapping to a full domain object may be partial. (2) The stream errors if the connection drops; wrap in `retryWhen` and resume from the last token. (3) Long-lived streams hold a cursor and use oplog; ensure the oplog window is large enough that a lagging consumer's token does not expire. (4) DELETE events give you the document `_id` only (no body, unless you configured pre-images in newer MongoDB). **When to use.** Event-driven side effects (cache invalidation, search-index sync, outbox-style propagation, notifications) where you need all operation types and reliable resumption. Use `@Tailable` only for the narrow case of tailing an append-only capped collection.

  • What deployment prerequisite do change streams have that @Tailable does not?
    Change streams require a replica set or sharded cluster because they read the oplog; a standalone mongod cannot serve them. @Tailable needs only a capped collection and works on a standalone.
  • Why request fullDocument = UPDATE_LOOKUP, and what is the trade-off?
    By default an UPDATE change event carries only the changed fields (updateDescription), not the whole document. UPDATE_LOOKUP fetches the current full document so you get a complete post-image, at the cost of an extra server-side lookup and possibly seeing a later state than the update itself.
  • What happens if your saved resume token is older than the oplog window?
    Resumption fails because the oplog no longer contains that position. You cannot continue gap-free and must fall back to a full resync of the affected data, then start a fresh stream.

saying these in an interview costs you the question

  • Claiming change streams work on a standalone mongod
  • Thinking change streams need a capped collection like @Tailable
  • Assuming UPDATE events always carry the full document by default
  • Believing the stream resumes automatically without persisting the token

context