<!-- llms-explorer concept facts · https://llms-explorer.com/tree/mongodb-kafka-connector/ · pack 2026-09-08 · ~3150 tokens -->

# mongodb-kafka-connector

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

Parent: [MongoDB Atlas Stream Processing](https://llms-explorer.com/tree/mongodb-atlas-stream-processing/) · 14 facets · 48 facts · page: https://llms-explorer.com/tree/mongodb-kafka-connector/

## 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 change events into Kafka for downstream consumers (analytics, search indexing, other microservices) without hand-rolling a change-stream-to-Kafka bridge. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Consuming Kafka events into MongoDB collections as a sink, for event-sourced or CQRS-style architectures where MongoDB is the read model. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Building a CDC pipeline that needs to survive connector downtime without losing events, when resume-token and oplog-window management matter. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Diagnosing why a Kafka Connect pipeline stopped processing MongoDB events, tracing it back to a specific source or sink connector error code. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Project ideas

- Build a CDC pipeline from MongoDB change streams to a downstream search index or cache, using the source connector with a topic naming convention like <prefix>.<db>.<collection>. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Implement an idempotent sink connector using ReplaceOneDefaultStrategy instead of the default InsertOneDefaultStrategy, so replayed Kafka messages don't fail on duplicate keys. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Build a resume-token-age monitor that alerts before the oplog window would invalidate the connector's stored token, distinct from ordinary Kafka consumer-lag monitoring. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Set up a Dead Letter Queue consumer that processes failed sink messages for alerting or manual replay instead of letting the connector halt on the first unprocessable message. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(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 — tasks.max should match the topic's partition count, or the extra tasks just contend with each other. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Not configuring a Dead Letter Queue in production, so the connector stops entirely on the first unprocessable message instead of routing it aside. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Using the default InsertOneDefaultStrategy for a pipeline that needs to be idempotent, causing duplicate-key failures on replay instead of using ReplaceOneDefaultStrategy or BulkWriteStrategy. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Watching Kafka consumer lag as the only health signal while ignoring resume-token age — consumer lag reflects the Kafka backlog, but token age reflects oplog risk, and an invalidated token forces a full resync. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Known issues

- If the connector's resume token expires because the oplog window was exceeded during downtime, it can't simply resume — it needs startup.mode=timestamp or a full re-snapshot with startup.mode=copy_existing. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- InvalidResumeToken errors can occur when a stored token is corrupted or comes from an incompatible MongoDB version, requiring the stored offset to be cleared and the connector restarted from scratch. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- DocumentValidationFailure (121) on the sink side means Kafka messages don't match the target collection's $jsonSchema validator — a schema mismatch that surfaces as a connector error rather than a clear validation message. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Sizing the oplog too small for the expected maintenance-window duration means any connector outage that exceeds the oplog window forces a full resync rather than a simple resume. — [source](https://llms-explorer.com/tree/mongodb-kafka-connector/) *(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)
