Schema evolution leaks the table

Raw CDC inherits source DDL, so an ALTER TABLE RENAME silently breaks downstream consumers — unless a schema registry enforces a compatibility mode that rejects the incompatible registration before any event ships.

Previously

Now the pipeline can faithfully replay history and stream live changes — but it is emitting raw row deltas with whatever column names the DBA picked, which means the next ALTER TABLE leaks straight into every consumer. So we need something that catches incompatible changes BEFORE the bytes leave the producer.

Scene 07

Schema evolution leaks the table

  1. Watch
  2. Try it
  3. Predict
  4. Capture
DDL paletteADDRENAMEDROPsource DBorders (id, user_id, amount, currency…)Schema Registryschema v1compatibility:BACKWARDFORWARDFULLNONEidle · waiting for DDLBACKWARD: new schema can read old dataDebeziumemits change events→ registers schemaKafka topicpayloads tagged v1old consumerv1✓ deserialize OKreading v1 cleanlynew consumerv2✓ deserialize OKreading v1 cleanly
What to watch for

The connector is streaming change events; the schema registry holds schema v1. Without a separate place to store schema versions, every consumer would have to redeploy on every DDL change. The component that does that storage is the schema registry, and the rules it enforces are called schema evolution. Watch what happens when the DBA adds a new column.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

SchemaRegistry.register
Gate a candidate schema by compatibility mode.
def register(subject, new_schema, mode):
old = versions[subject][-1]
if mode == NONE:
return accept(new_schema)
if mode == BACKWARD:
ok = can_read(reader=new_schema, data_of=old)
elif mode == FORWARD:
ok = can_read(reader=old, data_of=new_schema)
elif mode == FULL:
ok = (can_read(new_schema, old) and
can_read(old, new_schema))
return accept(new_schema) if ok else reject()
Connector.publishWithSchema
Register first, then embed the schema-id in the message.
def publish_with_schema(event, subject):
schema = avro_schema_of(event)
verdict = registry.register(subject, schema, MODE)
if verdict == REJECTED:
halt('cannot emit — incompatible schema')
schema_id = verdict.id
payload = magic_byte + schema_id + avro_encode(event)
kafka.produce(topic=subject, value=payload)
# consumer.deserialize() looks up schema by id
renameColumn
Avro rename = remove + add; aliases bridge the gap.
def rename_column(old_name, new_name, aliases_enabled):
new_schema = drop_field(current, old_name)
if aliases_enabled:
new_schema = add_field(new_schema, new_name,
aliases=[old_name])
else:
new_schema = add_field(new_schema, new_name)
# under BACKWARD, old reader looks for old_name:
# with alias -> resolves via aliases[] -> accept
# no alias -> field missing -> reject
return registry.register(subject, new_schema, MODE)

Where this sits in Build a CDC pipeline (Debezium + outbox)

Scene 07 of 12. Raw CDC inherits the source DDL — an ALTER TABLE RENAME COLUMN silently breaks downstream unless a Schema Registry enforces a compatibility mode that rejects the change at registration.

Up next. CDC of business tables couples consumers to internal table shape; the way out is to stop treating tables as the contract and start emitting domain events from a table that exists only to be published.

All 12 scenes in Build a CDC pipeline (Debezium + outbox) · Every curriculum

Built with Arqly
Every scene in Build a CDC pipeline (Debezium + outbox) builds on the one before it.All 12 Build a CDC pipeline (Debezium + outbox) scenes