Blog

Designing a Debezium CDC Flow That Survives Schema Changes

Table of Contents

Table of Contents

The source database had a harmless-looking ALTER TABLE. A column was added, the migration completed, and the application kept writing. The first sign of trouble appeared downstream: a Debezium task could no longer reconstruct the table state after a restart, a consumer rejected records with the changed schema, and a sink job began retrying because its destination mapping still expected the old shape.

This is an anonymous illustrative scenario, not a customer case. It shows how one DDL statement crosses several ownership boundaries. Debezium needs schema history to understand the source. Apache Kafka carries records and their serialized representation, while each consumer and sink decides what it can read and how it maps fields. Treating those as one schema problem leaves gaps between them.

The durable design is a chain of explicit contracts: protect the connector's schema history, define how the Topic carries versions, and make sink readers tolerate or deliberately reject the transition. The Kafka-compatible carrier can preserve the byte stream and Kafka semantics, but it does not replace the history store or decide whether a warehouse table should accept an added field.

A CDC schema change flow from source DDL through Debezium schema history and a Kafka Topic to a downstream sink

1The ALTER TABLE that traveled downstream

An ALTER TABLE is local to the database only until a change data capture pipeline turns it into an event contract. The source connector observes the DDL, updates its understanding of the table, and emits change records whose schema describes the row shape at that point in the log. A consumer that was deployed against the previous shape now has to decide whether it can deserialize, map, ignore, or quarantine the changed record.

The failure chain often looks longer than the original change:

  • The source connector restarts with incomplete or inaccessible schema history and cannot reconstruct the table definition needed for recovery.
  • The Kafka Topic contains records from more than one schema version, but consumers have no stated rule for reading both versions.
  • A sink maps fields into a fixed table or API payload and treats an added, removed, or renamed field as an unexpected input.

Each symptom points to a different control. Recreating the connector does not make a sink reader backward-compatible. A Schema Registry does not repair a deleted Debezium history Topic, and Kafka retention does not tell a destination how to apply a column rename. The pipeline needs a policy at each segment because the record changes meaning as it crosses the path.

Schema change resilience is not a single connector setting. It means recovering the source interpretation, carrying a mixed-version stream, and applying a known destination rule without changing data meaning.

2Schema history is the connector's reconstruction state

Debezium schema history records the source-side DDL knowledge that a connector uses to reconstruct table schemas while it reads the database log. For connectors such as MySQL, Debezium documents a Kafka-backed schema history store and the configuration that identifies it. The history is connector state. It is not a general-purpose contract for every consumer and it is not the same thing as the schema metadata carried with a change event. See the Debezium schema history documentation for the connector-specific boundary.

The distinction matters during recovery. A connector may need to interpret log records that refer to a table shape created by earlier DDL. With consistent history, it can rebuild the known schema sequence before resuming. If history is missing, points to the wrong Topic, or cannot be read with the configured identity, a restart can fail before a downstream application sees the record.

The most dangerous failure modes are mundane infrastructure mistakes:

  1. History lifecycle drift. An operator deletes the history Topic, applies incompatible retention settings, or restores the data plane without restoring history.
  2. Configuration mismatch. A replacement connector uses a different history Topic, bootstrap path, or credentials and sees an empty or unrelated stream.
  3. Permissions without proof. The connector can reach the data Topic but cannot read or write history. A healthy broker connection does not prove that every connector Topic works.
  4. Recovery without rehearsal. The team tests live flow but never stops and restarts the connector after a source DDL change.

The fix starts by making schema history a named dependency. Record its Topic, access policy, backup treatment, ownership, and recovery procedure alongside the connector configuration. Keep it separate from the business data contract so a sink team does not mistake reconstruction state for a public event interface. Verify connector-specific options in Debezium's schema storage configuration reference before writing the runbook.

One more boundary is worth stating plainly: preserving schema history does not make downstream evolution safe by itself. It only lets the source connector know what it is reading and how to describe the records it emits. The Topic and sink still need their own rules.

3Three segments, three evolution policies

The source, Topic, and sink see the same change from different angles. The source owns discovery and interpretation. The Topic owns the durable sequence of records and the serialization contract. The sink owns the destination mapping and its tolerance for mixed versions. A useful runbook gives each segment a policy, a test, and a clear failure action.

3.1Source connector: capture DDL and prove replay

The source policy should answer a narrow question: can the connector restart and reconstruct the table after the DDL has been committed? Start by deciding which source tables and schema changes are in scope. Then keep the Debezium history configuration under the same change control as the connector itself. A connector promotion that omits its history Topic, credentials, or access policy is an incomplete promotion.

The replay test should use an illustrative schema change in a non-production environment, emit records before and after it, stop the task, and start it from persisted state. Inspect the logs and records. The goal is to prove that the team can explain which schema applies at each boundary and recover when the task is moved or restarted.

Record the expected DDL scope, history access and recovery procedure, connector and converter versions, and the operator action when history cannot be read. That action should preserve evidence, repair the dependency, and resume from an approved offset or recovery point.

The source segment should fail loudly when reconstruction state is unavailable. Skipping it moves an early, diagnosable failure into a later sink incident.

3.2Topic: carry a readable contract across versions

The Topic is the seam between source interpretation and downstream interpretation. It should preserve the event envelope, keys, operation markers, source metadata, and serialized value rules that consumers rely on. During a transition, a CDC Topic can contain records from different schema versions, so “the Topic schema” is often a sequence of compatible shapes.

Not every evolution is safe. An added nullable field or reader-visible default can preserve old records. A removed or renamed field needs a reader migration, and a type change can fail even when its name stays the same.

Write those rules in a contract that downstream owners can test. Include the serialization format, compatibility direction, required fields, key behavior, tombstones, and the period during which legacy and current consumers may coexist. If a Schema Registry carries the contract, link its subject and compatibility policy in service documentation. If it is managed in code, make reader and fixture versions part of schema change approval. The wider ownership boundary is also useful in this Kafka data contract guide.

The Kafka layer has an important but limited role here. Apache Kafka documentation describes the Connect runtime, while the record contract belongs to the connector, converter, and consumers. Kafka retains the ordered byte stream according to Topic policy. It cannot infer that account_id is interchangeable with customer_id.

3.3Sink: map versions to destination behavior

The sink policy should be written in destination terms. A warehouse table may support an additive column but reject a type change. An API may ignore unknown fields but fail on a missing property. A search index may accept both fields while a materialized view needs a rebuild. “The sink supports schema evolution” is too broad to operate from.

For each destination, state how the sink handles an old record, a changed record, a missing field, an unexpected field, a renamed field, a type mismatch, a tombstone, and a replay. Decide which cases are retried, which are sent to a dead-letter path, and which stop the task for human repair. Kafka Connect's error handling configuration can control connector-level behavior, but it cannot define the business meaning of a column rename or choose the correct destination migration.

The sink deployment should follow the reader-first rule when the change is additive: make the destination and consumer able to accept the changed shape before the source begins emitting it. For a destructive change, keep the old field or mapping until all readers that depend on it have moved. The sink contract is complete only when the team can identify the exact record boundary that is safe to replay after a failed mapping.

Three evolution policies for the source connector, Kafka Topic, and downstream sink

4Downstream alignment: multi-version reads or a stop-write window?

Once each segment has a policy, the rollout still needs an alignment method. The choice is usually between keeping readers live across a mixed-version stream and creating a bounded window in which the source stops writing while every dependency moves together. Neither method removes the need for schema history. They reduce different kinds of uncertainty.

4.1Multi-version consumption keeps the stream moving

With multi-version consumption, deploy readers that can decode the old and current shapes before the source emits the change. The reader may use a default for an added field, accept both names during a rename, or route records by schema version to separate mapping code. The destination can then be prepared while Kafka continues to receive records.

This approach fits pipelines that need continuous capture and can carry compatibility logic for a transition. It also makes rollback more practical: an older reader can continue reading records only when the contract was designed for that direction. Compatibility must be tested against historical and current records, because a reader that handles the latest event can still fail when it catches up through older data.

The cost is a longer operating surface: more than one valid shape, more fixtures, and a removal decision. Give the compatibility path an owner and an exit condition showing that no deployment still depends on the legacy field.

4.2A stop-write window reduces concurrent states

With a stop-write window, pause source writes, drain the CDC stream to an agreed boundary, update the destination and sink mapping, verify connector and history state, and resume with the updated contract. State the boundary in offsets or source positions, not a vague maintenance promise. A schema compatibility gate can make that approval explicit before the pause begins.

This approach fits a destination that cannot read both versions and a business process that can tolerate a controlled pause. The state space is smaller because readers do not interpret legacy and current records at the same time. The cost is coordination: application, database, connector, and sink owners must agree on the pause, drain condition, and recovery action if one check fails.

Alignment methodUse whenMain controlMain riskProof before resume
Multi-version consumptionThe flow must keep capturing and readers can support both shapesReader compatibility, version-aware mapping, and a removal ownerLegacy and current behavior stays in production longer than plannedFixtures cover historical and changed records, and the sink accepts both paths
Stop-write windowThe destination cannot support both shapes and a coordinated pause is acceptableDrain boundary, ordered deployment, and explicit rollbackA partial change leaves a paused source or a mismatched sinkConnector restart, history read, destination mapping, and end-to-end write are checked

The table is a choice about where to hold complexity. Multi-version consumption holds it in readers and mappings; a stop-write window holds it in coordination and downtime. Choose the boundary you can observe and recover, then test it before the DDL reaches production.

5Where a Kafka-compatible carrier fits

The carrier matters because a CDC flow depends on Kafka Topic, Partition, Offset, and consumer behavior. A Kafka-compatible streaming platform such as AutoMQ can carry a Debezium and Kafka Connect path while the connector and consumer contracts remain explicit.

AutoMQ provides a Kafka-compatible data plane and a Shared Storage architecture. Its role here is to carry the Kafka stream through the broker storage layer, as described in the Kafka compatibility documentation and architecture overview. The storage layer is byte-stream neutral: it does not parse Debezium DDL, repair missing history, or decide whether a sink accepts a renamed field.

That boundary is useful during evaluation. Test history read and write, Topic serialization and retention, and sink mixed-version handling as separate acceptance tests. A compatible broker can remove a transport mismatch from the investigation, but it cannot remove ownership of schema history and downstream evolution.

6CDC resilience checklist

The fastest way to find a weak CDC design is to ask what happens after the DDL, a connector restart, and a sink replay. Use this checklist before approving a schema change:

  • Source history is owned. The Debezium history Topic, settings, credentials, access policy, and recovery treatment are documented with the connector.
  • DDL replay has evidence. A test proves that the connector can restart after the schema change and reconstruct the source schema.
  • Topic versions are intentional. Envelope, serialization, key behavior, tombstones, compatibility direction, and reader overlap are written down.
  • Sink behavior is explicit. Each destination has rules for field changes, retry, dead-letter, pause, and replay behavior.
  • Alignment is chosen. The team selects multi-version consumption or a stop-write window and names the proof boundary.
  • Recovery is separated by layer. History, deserialization, and destination mapping failures have different owners and runbooks.
  • The carrier is tested as a carrier. Compatibility, retention, ordering, Offset recovery, and connector Topics are checked without assigning schema decisions to the broker.

CDC resilience checklist covering source history, Topic contracts, sink behavior, alignment, and recovery

The anonymous ALTER TABLE scenario becomes manageable when each failure has a place to stop. The source connector proves it can reconstruct history. The Topic carries an intentional sequence of versions. The sink either reads both versions or waits for a controlled boundary. That is the difference between a schema change that travels through a pipeline and one that turns a local DDL into a chain of unrelated alerts.

If you are evaluating a Kafka-compatible carrier for this kind of CDC path, start with those acceptance tests and keep the schema responsibilities visible. Explore AutoMQ when you are ready to test the transport layer against your connector, Topic, and sink contracts.

Newsletter

Subscribe for the latest on cloud-native streaming data infrastructure, product launches, technical insights, and efficiency optimizations from the AutoMQ team.

Join developers worldwide who leverage AutoMQ's Apache 2.0 licensed platform to simplify streaming data infra. No spam, just actionable content.

I'm not a robot
reCAPTCHA

Never submit confidential or sensitive data (API keys, passwords, credit card numbers, or personal identification information) through this form.