Interview Setup
Interview Prompt
Design a CDC pipeline capturing changes from 100 source databases (5,000 tables) at 500K changes/sec (250 MB/sec), publishing to Kafka with schema evolution support and 150 TB retention.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Log based CDC (Debezium) or query based timestamp polling? | Log based capture observes committed row changes without repeated source queries. Polling can miss deletes and adds recurring query load. |
| One Kafka topic per table or per database? | 5,000 topics (one per table) enables independent consumer scaling but increases Kafka metadata overhead. |
| Schema evolution strategy: Avro with a registry or schemaless JSON? | An additive column change must not break downstream consumers. A schema registry provides explicit compatibility checks. |
| Delivery guarantee: at least once or exactly once to sinks? | At least once is simpler. Exactly once requires sink idempotency, version checks, or transactional integration. |
Scope
In scope
- End-to-end CDC pipeline flows and API contracts
- High-level architecture with major components
- Data model and storage choices with rationale
- Capacity estimation with shown math
- Primary failure modes and mitigations
Out of scope (state explicitly)
- Full analytics warehouse modeling (dbt/Star schema)
- Global exactly once side effects across all downstream consumers
- Building Kafka from scratch
Functional Requirements
Start by asking your interviewer which source databases you are capturing from and whether downstream consumers need exactly once side effects or can tolerate at least once delivery with replay safe sinks. Schema evolution, transaction boundaries, source log retention, and the initial snapshot strategy are common staff follow ups, so establish those assumptions before drawing Kafka.
- Capture every INSERT, UPDATE, DELETE from source databases in real time
- Propagate changes to downstream consumers with < 5 second latency
- Maintain strict ordering of changes per row/primary key
- Initial snapshot + ongoing incremental streaming
- Schema evolution: handle column adds, renames, type changes
- Multiple Sources: PostgreSQL, MySQL, MongoDB, SQL Server, Oracle
- Multiple Sinks: Kafka, Elasticsearch, S3/Parquet, Redis, ClickHouse
- Effectively Once Delivery via at least once capture plus replay safe, idempotent consumers and stable source identity
- Filtering and transformations: rename fields, mask PII, enrich
Non-Functional Requirements
Your interviewer will stress test source impact, per key ordering, transaction boundaries, replication slot WAL accumulation, snapshot recovery, and duplicate handling. Be explicit that ongoing CDC avoids polling while snapshots still read source tables and a stalled slot can consume source disk space.
- Low Latency: Database change reaches the downstream consumer in under 5 seconds (p99). Initial snapshots are tracked separately from steady state streaming latency.
- High Throughput: 100K+ changes/sec from a single source. A source that exceeds the throughput of one connector may require connector capable table splitting, source sharding, or a connector specific parallel capture strategy. This is connector and source engine dependent rather than a universal Debezium scaling rule.
- Low Source Impact: Ongoing CDC reads the transaction log instead of polling source rows. Initial and incremental snapshots still add source read load and must be scheduled and throttled appropriately.
- Effectively Once: At least once CDC combined with idempotent sink writes and replay safe versioning gives effectively once behavior per sink. External side effects such as payments or email still require destination specific idempotency or transactional controls.
- Durability: Capture every committed change while the replication mechanism remains valid and the required source log is retained. A source failure beyond the documented recovery assumptions is a separate failure domain.
- Schema Resilient: Compatible serialized schema changes can flow without breaking consumers when the compatibility policy and reader behavior support them. Breaking DDL changes require coordinated consumer rollout.
Capacity Estimations
Partition Kafka by primary key so per key ordering holds while the key mapping is stable. The baseline is 500K changes/sec across 100 source databases, so retention, topic metadata, hot tables, source log capacity, and connector CPU all matter to capacity planning.
| Metric | Calculation | Value |
|---|---|---|
| Source databases | Given (assumption documented in value) | 100 |
| Tables per database | Given (assumption documented in value) | 50 |
| Total tables | Given (assumption documented in value) | 5,000 |
| Changes / day (all sources) | 500K x 86400 | 43.2 Billion |
| Changes / sec (all sources) | 43.2B ÷ 86400 | 500K |
| Average changes / sec per source | 500K ÷ 100 | ~5K |
| Avg change event size | Given (typical workload assumption) | 500 bytes |
| Throughput | 500K x 500 bytes | 250 MB/sec |
| Kafka logical retention (7 days) | 250 MB/sec x 604800 sec | ~150 TB |
| Kafka replica storage (RF=3) | 150 TB x 3 | ~450 TB before filesystem and recovery headroom |
| Kafka topics (one per table) | Given (assumption documented in value) | 5,000 |
The ~150 TB figure is raw retained payload for the stated 250 MB/sec stream over 7 days. Production sizing also needs replication, indexes and segment overhead, compression variance, recovery headroom, and the configured disk utilization limit.
Architecture Diagram
Walk from left to right through the pipeline. Start with the source transaction log, move through the Debezium connector and Kafka partitions, then finish with schema management and downstream sinks. Ongoing capture reads the transaction log instead of polling source tables, while snapshots still perform controlled source reads.
Each PostgreSQL connector reads a replication slot on its source database. The slot keeps required WAL available while it remains valid, subject to the source disk budget. Partition Kafka by primary key so normal changes for one row land in one partition and per key ordering is preserved while the key mapping is stable. A primary key update is an important exception because the Kafka message key changes.
Downstream sinks should be designed for replay. Elasticsearch upserts and PostgreSQL upserts need stable source identity and, for mutable records, source version checks so an older replay cannot overwrite newer state. S3 and other object stores need deterministic file naming or committed file protocols when stronger guarantees are required. At least once delivery plus replay safe consumers provides effectively once behavior per sink without requiring a global distributed transaction.
In the room
Ask what happens when Debezium is down for six hours, because replication slot WAL accumulation can fill the source disk and take production offline. Mention slot-lag monitoring and max_slot_wal_keep_size before they probe it.
Component Deep Dives
How Log Based CDC Works
I start with log based CDC mechanics because the core design depends on reading the database transaction log without polling source rows. From there, the design moves through source connector differences, snapshot handoff, schema evolution, filtering, transaction boundaries, and the outbox pattern for domain events.
Log based CDC captures INSERT, UPDATE, and DELETE changes from the database transaction log without recurring source table polling. Lead with WAL and binlog reading before comparing it with polling or dual writes.
PostgreSQL Write-Ahead Log (WAL):
Every transaction writes WAL before data pages are modified
WAL is the source log from which logical replication decodes committed changes
Logical Replication (PostgreSQL 10+):
1. Create PUBLICATION: CREATE PUBLICATION my_pub FOR TABLE orders, users;
2. Debezium creates a REPLICATION SLOT
3. Debezium reads decoded WAL changes through the replication protocol
4. Each row change is decoded into: {table, op, before, after, source_metadata}
Replication slot behavior:
PostgreSQL retains WAL required by the slot while the slot remains valid and the required WAL has not been removed
This protects the connector from missing retained WAL during an outage, but it is not an unconditional data-loss guarantee if the slot is invalidated or required WAL is removed
Caution: A stalled connector can cause WAL growth and exhaust source disk space
Monitor: pg_replication_slots restart_lsn and confirmed_flush_lsn lag
Source impact:
Ongoing CDC reads the transaction log instead of repeatedly querying source tables
Initial snapshots do query source tables, so snapshot concurrency and timing must be controlled
Source CPU and I/O impact depend on transaction rate, row width, schema, decoding, and snapshot workload
Delete example:
UPDATE: before={...}, after={...}, op="u"
DELETE: before={...}, after=null, op="d"
A PostgreSQL table's REPLICA IDENTITY controls which old column values are available for UPDATE and DELETE events
DEFAULT normally provides the previous primary key values, whereas FULL provides the previous row image
A delete can be followed by a Kafka tombstone when the connector is configured to emit tombstonesSource Database Connector Differences
The overall pipeline is shared across database engines, but the change log and recovery mechanics differ by source.
| Source | Change Log | Key Recovery Concern |
|---|---|---|
| PostgreSQL | WAL + logical decoding | Replication slot retention and failover slot continuity |
| MySQL | Binlog | Binlog retention and GTID based recovery |
| MongoDB | Oplog / change stream | Resume token validity and oplog retention |
| SQL Server | Transaction log based CDC | Capture configuration and transaction log retention |
| Oracle | Redo based connector capture | Redo retention, privileges, and connector mode |
Kafka Connect Runtime and Task Scaling
Debezium normally runs inside a distributed Kafka Connect cluster. The REST API controls connector configuration, while workers execute connector tasks and store connector state in compacted internal Kafka topics. Scale workers for connector concurrency, but do not assume every source connector can split one database transaction log across arbitrary tasks without preserving source ordering.
Kafka Connect cluster:
Connect Worker 1 ─┐
Connect Worker 2 ─┼─> distributed worker pool
Connect Worker N ─┘
Connector:
PostgreSQL source
↓
Debezium task(s)
↓
Kafka topics
Internal state:
connect-configs → connector configuration
connect-offsets → source offsets / LSN positions
connect-status → connector and task status
Scaling rule:
Add workers for more connectors and task capacity
Increase source parallelism only when the connector preserves the required ordering semanticsFiltering, Masking, and Topic Routing
CDC platforms often need to remove sensitive fields, route tables into governed topic names, and apply lightweight record transformations before broad distribution. Keep transformations deterministic and make the source of truth remain the database log rather than the transformed stream.
Source change ↓ Debezium capture ↓ Column filtering / masking where policy allows ↓ Topic routing: cdc.<db>.<schema>.<table> ↓ Schema validation ↓ Kafka consumers / sinks Rule: Do not treat a transformed topic as the canonical source of database history Preserve source identity and ordering metadata for replay and reconciliation
Retries, Quarantine, and Dead Letter Queues
CDC pipelines need a clear boundary between transient failures and permanently invalid records. Retry network and downstream availability failures with bounded backoff. Send non retriable serialization or contract violations to a quarantine topic with the original source identity, partition, offset, and failure reason so operators can repair the consumer and replay the record later.
Processing:
Kafka record
↓
Validate / deserialize
↓
Retry transient failure with bounded backoff
↓
Success ───────────────> sink
│
└─ non retriable ──> quarantine / DLQ
DLQ record should retain:
source topic + partition + offset
table + primary key
source LSN / transaction metadata
schema version
failure reason
Rule:
Never silently drop a CDC record
Keep the original event replayable after the defect is fixedInitial Snapshot + Streaming
Every new connector must bootstrap existing rows before streaming. Explain the snapshot to streaming handoff and why incremental snapshots reduce the need for one long blocking operation on billion row tables.
Problem: When a CDC connector starts, the source table already has 100M rows. Debezium initial snapshot flow: 1. Start a transaction using the configured snapshot isolation level 2. Record the current database log position 3. Scan the captured tables and emit READ events 4. Finish the snapshot and record successful snapshot completion 5. Continue streaming changes from the recorded log position so updates committed during the snapshot are not omitted Important: Snapshot locking and isolation depend on connector settings The snapshot itself performs source table reads, so it adds bounded read load A connector that fails before snapshot completion can restart the snapshot For very large tables (billions of rows): Incremental snapshot: - Reads tables in chunks - Interleaves snapshot chunk processing with live streaming - Deduplicates snapshot rows against overlapping streamed changes - Runs without requiring one long table wide blocking operation
Change Event Format (Debezium Envelope)
The Debezium envelope is what every downstream consumer deserializes, so walk through before, after, op, and source metadata to demonstrate a firm grasp of the event contract.
{
"schema": {...},
"payload": {
"before": { "id": 42, "status": "pending", "amount": 99.99 },
"after": { "id": 42, "status": "completed", "amount": 99.99 },
"source": {
"connector": "postgresql",
"db": "ecommerce", "table": "orders",
"lsn": 987654321, "ts_ms": 1710403200000,
"txId": 12345678
},
"op": "u",
"transaction": {
"id": "12345678:987654000",
"total_order": 7,
"data_collection_order": 2
},
"ts_ms": 1710403200050
}
}Schema Evolution Without Downtime
Schema evolution is a staff level follow up that breaks naive pipelines, so explain Avro registry compatibility before describing what happens on ALTER TABLE ADD COLUMN vs DROP COLUMN.
Scenario: ALTER TABLE orders ADD COLUMN discount DECIMAL(10,2); What happens: 1. The source database applies the DDL 2. Debezium refreshes the captured table schema as needed and updates the emitted record schema 3. If a Schema Registry backed converter is configured, the new serialized schema is checked for compatibility 4. New events carry the evolved schema while older events remain readable under the configured compatibility policy 5. No connector restart is required for a compatible additive change Important limitation: PostgreSQL logical decoding does not emit generic DDL change events to consumers Schema metadata refresh and Schema Registry compatibility are separate concerns A breaking DDL change must be coordinated with downstream consumers Safer schema changes: Add a nullable or optional field: Usually compatible when the reader can handle its absence Remove a field: Compatibility depends on the configured schema evolution policy and consumer behavior Rename a field: Treat as a coordinated migration unless the downstream schema layer explicitly supports aliases Change a field type: Usually breaking unless the schema format and compatibility policy explicitly permit the widening Destructive changes should be deployed in phases so old and new consumers can coexist A compatibility check validates serialized schema evolution. It does not replace application level rollout testing
Backfill and Replay Strategy
Backfills are common when adding a new sink or repairing a projection. Use a historical Kafka offset or Debezium snapshot to rebuild the sink, isolate backfill consumers from live consumer groups, and make the sink replay safe before allowing the rebuilt projection to become authoritative.
Backfill: 1. Choose a source LSN, Kafka offset, or fresh snapshot boundary 2. Start an isolated consumer group 3. Rebuild the sink using idempotent upserts 4. Compare counts / checksums against the source 5. Replay the live delta after the chosen boundary 6. Cut over only after reconciliation passes
Exactly Once: Scope and Patterns
Effectively once sink behavior can be achieved without two phase commit by combining at least once delivery with replay safe, idempotent sinks. Use Kafka transactions when the processing topology benefits from Kafka atomicity, but do not extend a Kafka transaction claim to unrelated external systems.
Baseline source to Kafka semantics: Debezium provides at least once delivery by default Keying a record by primary key does NOT deduplicate updates because multiple valid updates can share the same key Kafka to sink: Elasticsearch: upsert by stable source key, guarded by source version or LSN when an older replay could otherwise overwrite newer state PostgreSQL: INSERT ... ON CONFLICT is idempotent when the sink applies the full current state or uses a source version guard for ordered updates S3 / Parquet: exactly once file output depends on the processing engine and sink implementation. Replay can create duplicate logical records unless the sink uses commit semantics or deterministic file generation Redis: SET of the same versioned value is idempotent, but counters and compare-and-set updates require explicit deduplication or version checks ClickHouse: use a schema and table engine that support replay-safe replacement or deduplication, and retain the source version or event identity needed to reconcile late or repeated changes Net effect: Use the term effectively once per sink when the sink is idempotent and replay safe Do not claim one global exactly once guarantee across Kafka, databases, object storage, Redis, and external systems without coordinated transactional integration
Transaction Boundaries and Primary Key Changes
Per key ordering is not the same as transaction ordering across multiple table topics. Enable Debezium transaction metadata when downstream logic needs transaction boundaries, and preserve the transaction identifier and event order when coordinating changes across topics. This is useful when a single source transaction updates several tables that feed independent Kafka topics. A primary key update also changes the Kafka message key. Debezium represents that change as a delete and create sequence, so consumers that depend on key based ordering must handle the transition explicitly.
Transaction A updates orders and order_items: BEGIN order row change → topic orders item row change → topic order_items END Transaction metadata: Enable the connector's transaction metadata option when needed transaction.id total_order data_collection_order Primary key change: old key → DELETE + tombstone new key → CREATE Debezium also carries headers that identify the new or old key for this transition Design rule: Per key ordering is preserved only while the key and partition mapping remain stable Cross table atomicity requires transaction aware consumer logic or an outbox/domain event
Outbox Pattern
CDC alone captures row changes rather than domain events. That makes the Transactional Outbox Pattern essential for emitting rich business events atomically with the database write. Learn more about event streams in Event Sourcing and CQRS.
Problem: The application needs to emit a DOMAIN EVENT rather than only a raw row change
CDC captures: "order.status changed from 'pending' to 'shipped'"
Downstream consumers instead need: "OrderShipped event with tracking_number and carrier"
Solution: Outbox Pattern
1. In the SAME database transaction:
UPDATE orders SET status='shipped' WHERE id=42;
INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ('order', '42', 'OrderShipped', '{"tracking":"1Z999"}');
2. Debezium CDC on the outbox table emits rich domain events to Kafka
3. Cleanup removes outbox rows only after CDC has safely advanced past them and the configured safety window has elapsed
Transactional guarantee:
The outbox row exists if and only if the business transaction commits
Emission to Kafka is asynchronous and normally at least once, so consumers must remain idempotentGovernance at 5,000 Topics
At 5,000 Kafka topics, metadata overhead and ACL sprawl become operational risks. Strict governance keeps the platform maintainable as teams self serve connectors.
One Kafka topic per table across 100 databases yields approximately 5,000 topics, which is manageable but metadata heavy. Governance rules enforce the naming convention cdc.{db}.{schema}.{table}, ensure automatic registration in the schema registry, and categorize topics by priority tier, such as P0 for orders and P3 for audit logs. Retention tiers range from 7 days of hot storage with cold S3 archival for P0 event topics down to 24 hours for lower priority event topics. Compaction is reserved for reconstructable state topics rather than historical audit streams. Access control lists isolate consumer teams so that service teams cannot inspect sensitive tables outside their domain. Monitor topic count growth, broker metadata size, and under replicated partitions. Cap new topics using a Debezium allowlist to prevent uncontrolled CREATE TABLE statements from spawning topics without review.
API Design
Kafka Connect Control Plane
Kafka Connect REST APIs are designed for operators, making them useful for showing how you provision connectors, inspect health, and pause or resume WAL consumption without restarting the pipeline.
POST /connectors
{
"name": "pg-ecommerce-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "pg-primary.internal",
"database.port": "5432",
"database.dbname": "ecommerce",
"database.user": "debezium",
"database.password": "<provided by secret manager>",
"topic.prefix": "pg-server1",
"plugin.name": "pgoutput",
"publication.name": "my_pub",
"table.include.list": "public.orders,public.users",
"column.exclude.list": "public.users.ssn",
"slot.name": "debezium_ecommerce",
"snapshot.mode": "initial"
}
}
GET /connectors
GET /connectors/{name}/status
POST /connectors/{name}/restart
PUT /connectors/{name}/pause
PUT /connectors/{name}/resumeCommon Error Responses
400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
202 Accepted: asynchronous job queued successfully, poll GET /jobs/{id} for completion status
408 Request Timeout: background job is still executing, continue polling status endpointData Model
Kafka Topic Layout
Baseline: one topic per captured table: pg-server1.public.orders -> partitioned by PK hash pg-server1.public.users -> partitioned by PK hash Key: serialized primary key, so the same row routes to the same partition while the key remains stable Value: Debezium change event envelope containing before, after, source, op, and optional transaction metadata Important edge case: Primary key update can change the message key Consumers must handle the old-key DELETE and new-key CREATE transition
Connect Offset Storage
Distributed Kafka Connect stores connector state in internal Kafka topics:
connect-offsets: compacted, replicated topic holding source connector offsets
connect-configs: compacted, highly replicated topic holding connector configuration
connect-status: compacted, replicated topic holding connector and task status
Example offset value:
Key: ["pg-ecommerce-connector", {"server": "pg-server1"}]
Value: {"lsn": 987654321, "txId": 12345678}
On restart: reads the saved connector offset and resumes from the corresponding WAL position when that WAL remains availableFault Tolerance
Source Log Capacity and Recovery Window
The source database has its own recovery window that is separate from Kafka retention. A seven day Kafka topic does not help if a stalled PostgreSQL logical slot loses the required WAL after the source exhausts its disk budget. Size the source log headroom for the expected outage and recovery time, then alert before the slot approaches its safe WAL limit.
Slot WAL Accumulation: The Primary Operational Risk
Problem: Debezium is down for 6 hours, causing PostgreSQL to retain 6 hours of WAL until disk fills up.
Critical: Production database runs out of disk space.
Mitigations:
1. Monitor: pg_replication_slots restart_lsn and confirmed_flush_lsn lag
Alert on lag that threatens the source disk budget or the recovery window
2. max_slot_wal_keep_size = 100GB (PostgreSQL 13+)
This limits how far a slot can fall behind before required WAL may be removed, which can invalidate the slot
If required WAL is removed and the slot becomes invalid, recover with a new consistent snapshot rather than blindly advancing to the current LSN
3. heartbeat.interval.ms: emit periodic heartbeat records so the connector can advance its acknowledged WAL position when captured tables are quiet
Optional heartbeat.action.query can update a heartbeat table when the source configuration requires it
4. PagerDuty alert if connector status != RUNNING for more than 5 minutes| Concern | Solution |
|---|---|
| Debezium crash | Restart the connector, read the saved source offset, and resume from that position when the required source log remains available. Reprocessing after the last checkpoint must be safe downstream |
| Source DB failure | Reconnect to the promoted primary using a valid recovery position. For PostgreSQL, synchronized failover slots can preserve logical slot state when configured. If required WAL is unavailable, take a new consistent snapshot and reconcile the sink. |
| Kafka broker failure | Replication factor 3 with min.insync.replicas=2 protects acknowledged events against a single broker failure when rack aware placement and clean leader election are maintained |
| Schema change (DDL) | Refresh captured schema metadata and validate serialized schema compatibility where Schema Registry is configured. Breaking DDL changes require coordinated consumer rollout |
| Duplicate events | Use source event identity such as table, primary key, and source log position or transaction metadata for downstream deduplication. Keep log compaction for reconstructable state topics, not as the duplicate defense for append only event streams. |
| Ordering violation | Partition by primary key so changes for the same stable key land on the same partition in strict order. Primary key changes require explicit delete and create handling. |
Additional Considerations
CDC vs Dual Writes vs Polling
| Approach | How | Pros | Cons | |---|---|---|---| | CDC (log based) | Read database WAL or binlog | Low ongoing source query load, captures committed DML changes, ordered per key | Operational complexity | | Dual writes | App writes to DB + Kafka | Simple | Race condition: DB succeeds, Kafka fails | | Outbox pattern | App writes to outbox, CDC reads outbox | Transactional consistency | Extra table, cleanup needed | | Polling (query) | SELECT WHERE updated_at > last_poll | Simple | Misses deletes, adds DB load, and introduces polling latency | CDC is usually a strong fit for real time synchronization, search index updates, cache invalidation, analytics pipelines, event driven microservices, and CQRS read model updates
Security and Data Governance
CDC carries production data into a shared event platform, so access control and data minimization are part of the architecture rather than an operational afterthought.
- Source access: Use dedicated least privilege replication users and grant only the schema, table, replication, and snapshot privileges required by the connector.
- Transport security: Require TLS for database, Kafka, Schema Registry, and Kafka Connect connections.
- PII controls: Exclude or mask sensitive columns before they reach broad consumer topics when the sink does not require them, and apply topic and consumer group ACLs by data domain.
- Secrets: Store database passwords, certificates, and Kafka credentials in a secret manager rather than connector source files.
- Auditability: Retain connector configuration history, schema versions, source LSNs, and downstream processing metrics so an event can be traced from the source transaction to each sink.
Multi Region CDC
Primary DB (us-east) -> Debezium -> Kafka -> MirrorMaker 2 -> Kafka (eu-west) Illustrative latency target: DB change to EU consumer ≈ 200ms under favorable network conditions Failover: promote the recovery region and redirect connectors/consumers only after source log and offset continuity is verified
Interview Walkthrough
- 25 minute cut
Skip the deeper architecture material unless the interview targets staff level.
- Debezium reads DB transaction log (8 min)
- One Kafka topic per table (9 min)
- At least once delivery (8 min)
- Explain CDC as reading the database transaction log such as PostgreSQL WAL or MySQL binlog rather than relying on polling or dual writes. For database internals, see Relational Database Engine (PostgreSQL).
- Cover effectively once delivery to downstream sinks through idempotent consumers, stable source identity, version checks, and offset based replay.
- Discuss schema evolution: additive changes are usually safe when readers tolerate missing fields, while column renames or deletions require coordinated multi phase migrations.
- Mention ordering guarantees per partition key when fanning out to downstream consumers like a Search Engine or cache invalidation layer.
- Cover lag monitoring and backfill strategy when adding a new downstream consumer.
- Common pitfall: application level dual writes to the database and search index guarantee data inconsistency when partial failures occur.
Engineering Trade-offs
Log Mining vs Query Polling vs Dual Write
Approach Latency DB Impact Completeness
-------------------------------------------------------------------------
Log Mining Low latency Minimal (read only) Captures committed DML changes,
(Debezium/WAL) Replication slot lag for captured tables, ordered
risk if consumer stalls per partition key
Query Polling Seconds to min Adds SELECT load Misses DELETEs,
(WHERE updated_at (poll interval) on hot tables misses rows updated
> last_seen) twice within gap
Dual Write Milliseconds App-side complexity Inconsistent if app
(App -> DB + Kafka) Two network calls crashes after DB
write but before KafkaExactly Once and Replay Safety: When It Matters
Exactly once behavior matters when duplicate side effects have business impact, such as financial ledger entries or inventory decrements. Search index updates, aggregate analytics, and cache warming usually prefer replay safe idempotency because a global transaction across every system is expensive. The correct guarantee is therefore stated per processing topology and sink, not as one universal property of the CDC platform.
Common mitigation techniques include stable source identity, source version or LSN checks, idempotent upserts, deduplication records keyed by source event identity, and Kafka transactional semantics where the processing topology supports them. For external side effects such as email or payments, use a provider idempotency key or transactional integration where available.
Schema Evolution in a Live Pipeline
A safe multi phase rollout uses four ordered steps. First, add the column as nullable with no default in the source database. Second, deploy consumer code that can handle both schema variants. Third, deploy writers that populate the new field. Fourth, apply the NOT NULL constraint after backfilling existing rows. Schema validation failures should be isolated and replayable rather than silently dropped.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.