MongoDB to Kafka: How to Move Your Data

Move MongoDB into Kafka with Airbyte. Why the oplog and topic retention windows both matter, and how the schema mode decides what consumers receive.

Summarize with AI:

Moving MongoDB onto Kafka publishes application changes where other services can react to them. A document store holding orders or user state is often the system of record, and several teams wanting to know when something changes should not all be reading your database.

This guide covers the managed path with Airbyte. Two things shape the build: there are two retention windows in this pipeline rather than one, and a schema decision made at setup determines what every consumer receives.

MongoDB to Kafka at a glance:

CapabilitySupportedWhat it means for this pipeline
Change captureChange streamsSo inserts, updates and deletes all reach the topic
Replica setRequiredChange streams do not exist on a standalone mongod
Oplog retentionYour constraintSync less often than it holds and the offset is lost
Schema modeTwo optionsTyped fields from sampling, or the whole document
Discovery sampling10,000 by defaultConfigurable, and sparse fields can be missed

Why move data from MongoDB to Kafka?

Two situations account for most of these pipelines.

The first is decoupling. Several services wanting to know when a document changes can subscribe to a topic instead of polling your database or having application code call each of them, which is the difference between an architecture and an accumulation of integrations.

The second is feeding an existing event-driven platform where Kafka is already the backbone. If only one system needs this and the goal is analysis rather than reaction, MongoDB to Databricks is simpler to build and far easier to query afterwards.

What do you need before you start?

Four things, and the second is where most trouble with this connector originates:

A replica set, not a standalone server. Change streams exist only on a replica set or sharded cluster, so a standalone mongod cannot be used. The MongoDB source documentation covers the requirements and the privileges.

An oplog sized for your sync schedule. Airbyte recommends it hold at least twenty-four hours of changes and suggests a week for comfort, because a sync that falls outside that window cannot resume.

Topics created in advance. The destination writes to topics that already exist, and one per collection is almost always better than combining them.

A decision about schema enforcement. The connector enforces a schema by default, sampling documents to decide which fields exist, and the alternative passes each document through whole. Your consumers live with whichever you pick.

If your cluster restricts inbound traffic by IP, add the Airbyte Cloud IP addresses to the allow list before you begin.

How do you build a MongoDB to Kafka pipeline in Airbyte?

Step 1: Check your oplog retention before anything else

Find out how many hours of changes your oplog currently holds, and set it to a week if you have any choice in the matter. This is the single most common cause of trouble with this connector, and the failure is not subtle: the connector reports that the saved offset is no longer valid and asks you to reset the connection, which means a full resync of every collection you were being careful about.

Step 2: Configure the MongoDB source

Click Sources in the left navigation, then New Source, and select MongoDB, following adding a source. Supply the connection string, database and credentials, then choose your schema mode and sample size. Pointing discovery at a secondary is worth doing if the cluster is busy.

Step 3: Configure the Kafka destination

Click Destinations, then New Destination, and select Kafka, following adding a destination. Supply the bootstrap servers, security protocol and topic configuration. Messages are JSON wrapping each record with its identifier and stream name, so consumers read through an envelope.

Step 4: Create the connection and sync frequently

Click Connections, then New connection, select your collections and a sync mode. Frequency matters more here than usual, because the gap between syncs has to stay comfortably inside your oplog window. Alert on failure, since a pipeline stopped for a weekend may not be able to resume.

Then tell your consumers what shape to expect, because that is settled by a configuration choice they cannot see.

Why are there two retention windows here?

Because both ends of this pipeline forget things on a schedule, which is unusual. MongoDB's change streams read from the replica set oplog, a fixed-size log that discards old entries, and Kafka topics have their own retention policy. Data has to cross from one window into the other before either closes.

The oplog is the one that produces the worse failure. If the pipeline pauses longer than the oplog holds, the position it saved no longer exists, and the connector says so plainly: the saved offset is not valid, reset the connection and increase oplog retention or sync frequency. Resetting means resyncing everything, which on a large collection is an afternoon rather than a moment.

Both remedies are within your control and neither is clever. Extend the oplog so it holds a week, sync frequently enough that a weekend outage stays inside it, and alert on failure so nobody discovers this on Monday. Then check that your Kafka retention gives consumers time to read what arrives, since an event that expires before a consumer catches up is just as lost.

What shape do your consumers actually receive?

Whichever one you chose at setup, and they have no way of knowing which. With schema enforcement on, the connector samples documents during discovery and produces a fixed set of typed fields, so consumers get something predictable. With it off, each document passes through whole, so consumers get everything and no guarantees.

Neither is obviously right, and the trade is worth naming. Enforcement gives a stable contract and can silently omit fields that were too rare to appear in the sample, which defaults to ten thousand documents. A field present on two percent of documents may simply not be in the schema, and a consumer depending on it will never see it arrive.

For a bus feeding services rather than a warehouse feeding reports, the calculation often favours passing documents through whole, since a consumer can ignore fields it does not want and cannot invent ones that never arrived. Whichever you choose, write it down beside the topic name, because a downstream team debugging a missing field has no way to discover that a sampling decision made months ago is the reason.

Frequently asked questions

The sync says the saved offset is not valid.

The pipeline fell outside your oplog retention window. Reset the connection, then increase oplog retention or sync more frequently so it does not recur.

Do I need a replica set?

Yes. Change streams exist only on a replica set or sharded cluster, so a standalone mongod cannot be used for change capture.

A field is missing from the messages.

Schema discovery probably did not see it in the sampled documents. Raise the sample size and rerun discovery, or use the schemaless mode so documents pass through whole.

Do deletions reach the topic?

Yes, because change streams capture inserts, updates and deletes, which is one of the better reasons to use this connector rather than polling.

Can I do this without writing code?

The pipeline, yes, though sizing the oplog is database administration. The consumers are yours and their expectations are the part worth documenting.

Get your MongoDB data into Kafka

Size the oplog for a week before you build anything, because falling outside it costs a full resync and the error tells you so explicitly. Confirm you have a replica set. Sync frequently and alert on failure, since a weekend outage is the usual way this breaks. Then decide the schema mode deliberately and record it beside the topic, because your consumers receive whichever shape you picked and cannot tell which.

Airbyte's connector catalog includes 600+ pre-built connectors, so application changes can reach every service that cares about them. For a relational source onto the same bus, see MySQL to Kafka, and for the same source into a warehouse, MongoDB to Snowflake.

Start syncing now →

Integrate with 700+ apps using Airbyte

Move data from 700+ sources into warehouses, lakes, and beyond. Set up pipelines in minutes with pre-built connectors and the Connector Builder.