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

# kafka-source-connector

> 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/kafka-source-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 topics for downstream consumers, instead of polling MongoDB directly from every consuming service. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Building a near-real-time CDC pipeline where consumer lag and latency matter, tuning poll.await.time.ms and poll.max.batch.size to trade throughput for lower latency. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Recovering a source connector after extended downtime caused the oplog window to be exceeded and the stored resume token to expire. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Deciding between resuming from a timestamp versus a full re-snapshot (copy_existing) after a ChangeStreamHistoryLost error. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Project ideas

- Stand up a basic MongoDB change-stream source connector publishing to a topic named mongo.<db>.<collection>, then simulate an outage long enough to exceed the oplog window and practice recovering with startup.mode=timestamp. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Build a low-latency CDC pipeline by setting poll.await.time.ms to 100 and poll.max.batch.size to a small value between 100 and 500, then measure the consumer-lag improvement against default settings. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Monitor both Kafka consumer lag and resume token age side by side on a running source connector, since they signal different kinds of risk. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Trigger and resolve an InvalidResumeToken error by clearing the stored offset and restarting the connector against a fresh resume point. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Common mistakes

- Monitoring only Kafka consumer lag and not resume token age — consumer lag tells you about the Kafka-side backlog, but resume token age tells you whether the oplog window is at risk of invalidating the token entirely. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Not sizing the oplog large enough to cover expected connector downtime or maintenance windows, forcing a full resync when a routine outage runs long. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Treating ChangeStreamHistoryLost (error code 286) as unrecoverable instead of resolving it with startup.mode=timestamp or a copy_existing re-snapshot. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Leaving the connector on default poll settings for a use case that actually needs low latency, instead of deliberately tuning poll.await.time.ms and poll.max.batch.size. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*

## Known issues

- If the resume token expires because the oplog window was exceeded during connector downtime, the connector can't just pick back up — it needs either a timestamp-based restart or a full copy_existing re-snapshot. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- InvalidResumeToken can occur not just from oplog truncation but from a resume token that's corrupted or from an incompatible MongoDB version, and the only resolution is clearing the stored offset and restarting. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- Change events publish to a fixed topic-naming convention (<prefix>.<db>.<collection>), which constrains downstream topic routing and consumer subscription design. — [source](https://llms-explorer.com/tree/kafka-source-connector/) *(AI-suggested, synthesized from this pack's existing facts — not extracted from a source document.)*
- The source connector persists its resume token in a Kafka Connect offsets topic, so recovery behavior is tied to that offsets topic's own durability and retention, not just to MongoDB's oplog. — [source](https://llms-explorer.com/tree/kafka-source-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)
