14 Nov 2025 · 8 min read · Vivek Amethiya
Change data capture with Debezium and Kafka: SQL Server to MongoDB
How log-based CDC works, how Debezium turns SQL Server changes into Kafka events, and the patterns that keep a downstream MongoDB copy correct.
- Kafka
- Debezium
- CDC
- SQL Server
- MongoDB
- Node.js
Two databases that must hold the same truth are one of the most common problems in backend work. On the eHUB brokerage platform I led a pipeline that moved investor account data from Microsoft SQL Server into MongoDB every day, with full consistency. The approach that made it reliable was change data capture (CDC) with Debezium and Apache Kafka.
This article explains how that stack works, end to end, and the patterns that keep the copy correct in production. Table and database names below are illustrative.
Three ways to keep two databases in sync
Dual writes. The application writes to both databases. It looks simple until one write succeeds and the other fails. There is no transaction that spans SQL Server and MongoDB, so the two drift apart and nobody notices until a customer does.
Batch ETL. A scheduled job queries the source and copies rows over. It is easy to reason about, but every run re-reads data that has not changed, it puts load on the production database, and it struggles with deletes: a row that no longer exists is not returned by any query.
Change data capture. The pipeline reads the database's own record of what changed, in the order it was committed, and replays it downstream. Only changes move, deletes are included, and the source application does not change at all.
For a brokerage, where a missed delete or an out-of-order update means an account shows the wrong state, CDC was the clear choice.
What log-based CDC actually reads
Every relational database keeps a transaction log so it can recover after a crash. Each committed insert, update and delete is written there before it is applied to the tables. Log-based CDC reads that log instead of querying tables.
SQL Server packages this as a built-in feature. Once CDC is enabled, a capture job run by SQL Server Agent reads the transaction log and writes every change for the tracked tables into change tables in a cdc schema. Debezium reads those change tables.
Enabling it is two calls, run by a member of sysadmin and db_owner respectively:
-- Once per database
EXEC sys.sp_cdc_enable_db;
-- Once per table you want to stream
EXEC sys.sp_cdc_enable_table
@source_schema = N'dbo',
@source_name = N'accounts',
@role_name = NULL,
@supports_net_changes = 0;SQL Server Agent must be running, because the capture and cleanup jobs are Agent jobs. If Agent is stopped, the change tables stop filling and the pipeline silently stalls.
The pipeline
Debezium runs as a source connector inside Kafka Connect. You register it by posting a JSON configuration to the Connect REST API:
{
"name": "brokerage-sqlserver",
"config": {
"connector.class": "io.debezium.connector.sqlserver.SqlServerConnector",
"database.hostname": "mssql.internal",
"database.port": "1433",
"database.user": "debezium",
"database.password": "${file:/opt/secrets/mssql.properties:password}",
"database.names": "brokerage",
"topic.prefix": "ehub",
"table.include.list": "dbo.accounts,dbo.positions",
"snapshot.mode": "initial",
"decimal.handling.mode": "string",
"tombstones.on.delete": "true",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schema-history.ehub"
}
}A few settings matter more than they look:
topic.prefixnames the topics. Each table gets its own:ehub.brokerage.dbo.accounts.snapshot.mode: initialcopies the existing rows once, then switches to streaming changes. Without it, the target starts empty.decimal.handling.mode: stringsends money as text such as"2500.00". The default encodes decimals as bytes, and thedoubleoption loses precision. For balances, use strings.schema.history.internal.*is where Debezium remembers table structures, so it can decode older changes after a restart.- The password uses a Kafka Connect config provider, so the secret is read from a file and never stored in the connector config.
What a change event looks like
Every message has a key (the row's primary key) and a value called the envelope:
{
"before": { "account_id": 1042, "status": "PENDING", "cash_balance": "2500.00" },
"after": { "account_id": 1042, "status": "ACTIVE", "cash_balance": "2500.00" },
"source": {
"connector": "sqlserver",
"db": "brokerage",
"schema": "dbo",
"table": "accounts",
"commit_lsn": "0000002a:00000c38:0005"
},
"op": "u",
"ts_ms": 1727856000123
}opis the operation:ccreate,uupdate,ddelete,ra row read during the initial snapshot.beforeandafterare the row before and after the change. A create has nobefore, a delete has noafter.sourcesays exactly where the change came from, including its position in the log (commit_lsn).
With tombstones.on.delete enabled, a delete is followed by a second message with the same key and a null value. It exists for Kafka log compaction, so consumers must skip it.
Applying changes safely in Node.js
Kafka delivers messages at least once. After a crash or a rebalance, a consumer may see the same event again. The rule that follows is simple: every write must be idempotent. Applying the same event twice must leave the same result as applying it once.
Keying by primary key gives the other guarantee we need: all changes to one row land in the same partition, in commit order. So an upsert keyed on the primary key is both idempotent and ordered:
import { Kafka } from 'kafkajs';
import { MongoClient } from 'mongodb';
interface AccountRow {
account_id: number;
status: string;
cash_balance: string;
}
interface ChangeEvent {
before: AccountRow | null;
after: AccountRow | null;
op: 'c' | 'u' | 'd' | 'r';
ts_ms: number;
}
const kafka = new Kafka({ clientId: 'accounts-sync', brokers: ['kafka:9092'] });
const consumer = kafka.consumer({ groupId: 'accounts-sync' });
const mongo = new MongoClient(process.env.MONGO_URL!);
const accounts = mongo.db('ehub').collection('accounts');
await mongo.connect();
await consumer.connect();
await consumer.subscribe({ topic: 'ehub.brokerage.dbo.accounts', fromBeginning: true });
await consumer.run({
eachMessage: async ({ message }) => {
if (!message.value) return; // tombstone: the delete event before it did the work
const event = JSON.parse(message.value.toString()) as ChangeEvent;
if (event.op === 'd' && event.before) {
await accounts.deleteOne({ _id: event.before.account_id });
return;
}
if (event.after) {
const { account_id, status, cash_balance } = event.after;
await accounts.updateOne(
{ _id: account_id },
{ $set: { status, cashBalance: cash_balance, syncedAt: new Date(event.ts_ms) } },
{ upsert: true },
);
}
},
});The offset is committed only after eachMessage resolves, so a crash mid-write means the event is delivered again and the upsert simply repeats.
Keep slow work off the stream
Not every reaction to a change belongs in the consumer. If an account turning ACTIVE should trigger a statement rebuild or a file transfer, doing it inline would hold up every change behind it. On eHUB that kind of work went to BullMQ, a Redis-backed job queue:
import { Queue } from 'bullmq';
const jobs = new Queue('account-jobs', { connection: { host: 'redis', port: 6379 } });
// Inside eachMessage, after the upsert:
if (event.op === 'u' && event.before?.status !== 'ACTIVE' && event.after?.status === 'ACTIVE') {
await jobs.add(
'rebuild-statement',
{ accountId: event.after.account_id },
{
jobId: `statement-${event.after.account_id}-${event.ts_ms}`, // duplicates are ignored
attempts: 5,
backoff: { type: 'exponential', delay: 2000 },
},
);
}The jobId is derived from the event, so a replayed event cannot enqueue the same job twice. Retries with exponential backoff absorb a flaky downstream system without blocking the stream.
What to watch in production
The CDC retention window. SQL Server's cleanup job deletes change data after three days by default. If the connector is down longer than that, changes are gone and you need a fresh snapshot. Alert on connector health well inside that window, and tune retention with sys.sp_cdc_change_job if you need more headroom.
Schema changes. A SQL Server capture instance does not pick up a new column on its own. Adding a column means creating a second capture instance for the table; Debezium then switches to it. Plan schema changes with the pipeline in mind, not after.
Ordering has limits. Order is guaranteed per key, not across tables. An order line can arrive before its order. Design consumers to tolerate that, for example by upserting a parent placeholder or retrying the child later.
The initial snapshot. Snapshotting large tables takes time and, under the default isolation level, can hold locks. If the database allows it, snapshot.isolation.mode: snapshot avoids blocking writers while the snapshot runs.
Poison messages. One event that always fails to apply will stop a partition. Route repeated failures to a dead-letter topic with the error attached, and alert on it.
Lag. Consumer lag per partition is the single best health signal. Rising lag means the target is falling behind the source, long before anyone sees stale data.
When CDC is not worth it
If the data is small, changes rarely, and a nightly copy is acceptable, a scheduled export is simpler to run. CDC earns its operational cost when you need every change, including deletes, delivered in order with low delay.
Summary
- Log-based CDC replays the database's own record of changes, so nothing is missed and the source application stays untouched.
- Debezium turns SQL Server change tables into ordered Kafka events, one topic per table, keyed by primary key.
- Consumers must be idempotent: upsert by key, skip tombstones, commit after writing.
- Push slow side-effects to a job queue with deterministic job ids.
- Monitor retention, schema changes, poison messages and lag.
You can see how this fits into the full eHUB architecture in the case study.