<!-- llms-explorer concept facts · https://llms-explorer.com/tree/cdc-patterns/ · pack 2026-09-08 · ~3094 tokens -->

# CDC-patterns

> The MongoDB Connector for Apache Kafka is a Kafka Connect plugin that bridges MongoDB and Kafka in both directions:

Parent: [mongodb-kafka-connector](https://llms-explorer.com/tree/mongodb-kafka-connector/) · 14 facets · 48 facts · page: https://llms-explorer.com/tree/cdc-patterns/

## Overview

- The MongoDB Connector for Apache Kafka is a Kafka Connect plugin that bridges MongoDB and Kafka in both directions: — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#overview)
  - Source connector: MongoDB change streams → Kafka topics (CDC pipeline) — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#overview)
  - Sink connector: Kafka topics → MongoDB collections (event consumer) — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#overview)
- Supports Confluent Platform, Confluent Cloud, Amazon MSK, and self-managed Kafka. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#overview)

## Basic Change Stream Source

- This publishes change events to topic mongo.mydb.orders (format: <prefix>.<db>.<collection>). — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#basic-change-stream-source)

## Resume Token Persistence

- The connector automatically persists the resume token in a Kafka Connect offsets topic. On restart, it resumes from the saved token. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#resume-token-persistence)
- If the token expires (oplog window exceeded during connector downtime): — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#resume-token-persistence)

## Dead Letter Queue (DLQ)

- Configure DLQ to route failed messages instead of stopping the connector: — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#dead-letter-queue-dlq)
- DLQ messages include headers with error context. Process DLQ messages with a separate consumer for alerting or manual replay. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#dead-letter-queue-dlq)

## Source Connector Errors

- ChangeStreamHistoryLost (error code 286): — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#source-connector-errors)
- The oplog has been truncated past the resume token. Resolution: — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#source-connector-errors)
  - Set startup.mode=timestamp to start from a recent time — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#source-connector-errors)
  - Or re-snapshot with startup.mode=copy_existing — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#source-connector-errors)
  - Increase oplog size to prevent future occurrences — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#source-connector-errors)
- InvalidResumeToken: Resume token is corrupted or from an incompatible MongoDB version. Resolution: clear stored offset and restart connector. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#source-connector-errors)

## Sink Connector Errors

- DuplicateKey (11000): Configure ReplaceOneDefaultStrategy instead of InsertOneDefaultStrategy to make sink idempotent. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#sink-connector-errors)
- DocumentValidationFailure (121): Kafka messages don't match MongoDB $jsonSchema validator. Check message schema vs collection validator. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#sink-connector-errors)

## Sink Connector

- Worker parallelism: Set tasks.max equal to the number of Kafka partitions for the topic. — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#sink-connector)

## CDC Pipeline Pattern: MongoDB → Kafka → Downstream

- For near-real-time with low latency: — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#cdc-pipeline-pattern-mongodb-kafka-downstream)
  - Use poll.await.time.ms: 100 (shorter poll interval) — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#cdc-pipeline-pattern-mongodb-kafka-downstream)
  - Monitor consumer lag on the Kafka topic — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#cdc-pipeline-pattern-mongodb-kafka-downstream)
  - Keep poll.max.batch.size small (100-500) for lower latency at cost of throughput — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#cdc-pipeline-pattern-mongodb-kafka-downstream)

## Anti-Patterns

- Single-partition topics with multiple sink tasks: Multiple sink tasks on a single partition = contention; match tasks.max to partition count — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#anti-patterns)
- Not configuring DLQ: Connector stops on first unprocessable message; always configure DLQ in production — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#anti-patterns)
- InsertOneDefaultStrategy for idempotent pipelines: Insert fails on duplicate; use ReplaceOneDefaultStrategy or BulkWriteStrategy for idempotent sinks — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#anti-patterns)
- Monitoring consumer lag but not resume token age: Consumer lag tells you about Kafka backlog; resume token age tells you about oplog risk (if token becomes invalid = full resync) — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#anti-patterns)
- Not increasing oplog for high-volume CDC: Connector outage exceeding the oplog window = full resync required; size oplog to cover expected maintenance windows — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#anti-patterns)

## References

- MongoDB Kafka Connector Documentation — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#references)
- Source Connector Configuration — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#references)
- Sink Connector Configuration — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#references)
- Write Model Strategies — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#references)
- Kafka Connector GitHub — [source](https://llms-explorer.com/sources/mdb-context-hub/mongodb-kafka-connector/#references)

## Where this helps

- Streaming MongoDB changes to downstream systems in near-real-time via Kafka, a CDC pipeline, rather than polling the database on an interval. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- A Kafka Connect worker restarts and resume-token persistence needs to be understood so the source connector picks back up from where it left off instead of replaying or dropping events. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- A sink connector is failing on duplicate-key errors and the right write strategy, ReplaceOneDefaultStrategy versus InsertOneDefaultStrategy, needs to be chosen to make it idempotent. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Sizing sink connector parallelism and matching tasks.max to the number of Kafka partitions to avoid contention. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Project ideas

- Stand up a MongoDB change-stream source connector publishing to a topic named by the prefix.db.collection convention, then trace a document update through to the Kafka topic. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Configure a dead-letter queue on a sink connector and deliberately send a malformed message to confirm it routes to the DLQ with error-context headers instead of stopping the connector. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Build a resume-token recovery path that detects a ChangeStreamHistoryLost condition and falls back to startup.mode=timestamp to restart from a recent point rather than failing permanently. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Tune a low-latency source pipeline by lowering poll.await.time.ms and monitoring consumer lag on the downstream topic to confirm the change actually reduces end-to-end latency. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Common mistakes

- Running multiple sink tasks against a single-partition topic, which creates contention instead of parallelism; tasks.max should match the partition count. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Not configuring a dead-letter queue, so the connector stops entirely on the first unprocessable message instead of routing it aside and continuing. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Using InsertOneDefaultStrategy for a pipeline that needs idempotent replays, which fails on duplicate keys instead of the intended replace or upsert behavior. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Letting the oplog window get exceeded during connector downtime, which invalidates the persisted resume token and forces a restart from a timestamp rather than a clean resume. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Known issues

- A resume token can expire if the connector is down longer than the oplog retention window, forcing a startup.mode=timestamp restart and risking missed or re-processed events around the gap. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- ChangeStreamHistoryLost, error code 286, means the oplog has already been truncated past the resume token; by the time this error surfaces, the original position is unrecoverable and the fix is inherently lossy. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- DocumentValidationFailure, error 121, on the sink side means Kafka messages don't match the target collection's $jsonSchema validator; schema drift between producer and consumer isn't caught until write time. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Sink connector idempotency depends entirely on choosing the right write strategy per use case; the default strategy isn't automatically safe for every workload, so this has to be a deliberate configuration choice. — [source](https://llms-explorer.com/tree/cdc-patterns/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Context files

- [mongodb-kafka-connector](https://llms-explorer.com/downloads/sources/mdb-context-hub/mongodb-kafka-connector.md)
