Implementierung von Change-Data-Capture-Pipelines für Datenbanken
Baue CDC-Pipelines, die Datenbankänderungen in Echtzeit streamen: logbasierte Erfassung mit Debezium, Event-Format, Schema-Evolution, Exactly-once.

Eine Datenbank alle paar Sekunden per Polling auf Änderungen abzufragen, verschwendet Ressourcen und erzeugt Latenz. Change Data Capture (CDC) liest das eigene Transaktionslog der Datenbank — das Write-Ahead-Log (WAL) in PostgreSQL, das Binlog in MySQL — und streamt jedes Insert, Update und Delete als Event. Damit erhältst du Echtzeit-Datenreplikation, ohne die Performance der Quelldatenbank zu beeinträchtigen.
CDC ist das Rückgrat moderner Datenarchitekturen: Es treibt Search-Index-Updates, Cache-Invalidierung, Analytics-Pipelines und die datenbankübergreifende Synchronisation zwischen Services an, ohne Producer an Consumer zu koppeln.
Logbasiertes CDC mit Debezium
Debezium liest das Transaktionslog der Datenbank und veröffentlicht Change-Events in Kafka. Die Quelldatenbank braucht keine Anpassung — keine Trigger, keine Polling-Queries, kein Event-Publishing auf Anwendungsebene.
# ❌ Polling-based approach — high latency, wasteful queries
# SELECT * FROM orders
# WHERE updated_at > :last_check_time
# Run every 5 seconds across all tables
# Misses deletes, creates load on source database// ✅ Debezium connector configuration for PostgreSQL CDC
interface DebeziumConnectorConfig {
name: string;
config: {
"connector.class": string;
"database.hostname": string;
"database.port": number;
"database.user": string;
"database.dbname": string;
"database.server.name": string;
"plugin.name": string;
"slot.name": string;
"publication.name": string;
"table.include.list": string;
"transforms": string;
"transforms.unwrap.type": string;
"transforms.unwrap.drop.tombstones": string;
"key.converter": string;
"value.converter": string;
"value.converter.schemas.enable": string;
};
}
const ordersCdcConnector: DebeziumConnectorConfig = {
name: "orders-cdc-connector",
config: {
"connector.class":
"io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "db.internal",
"database.port": 5432,
"database.user": "cdc_reader",
"database.dbname": "production",
"database.server.name": "prod-orders",
"plugin.name": "pgoutput",
"slot.name": "orders_slot",
"publication.name": "orders_publication",
"table.include.list": "public.orders,public.order_items",
"transforms": "unwrap",
"transforms.unwrap.type":
"io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"key.converter":
"org.apache.kafka.connect.json.JsonConverter",
"value.converter":
"org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
},
};Das pgoutput-Plugin liest den logischen Replikationsstream von PostgreSQL. Die ExtractNewRecordState-Transformation flacht das Envelope-Format von Debezium auf den Nachher-Zustand jeder Zeile ab und vereinfacht so die nachgelagerte Verarbeitung.
Change-Events verarbeiten
Jedes CDC-Event enthält den Operationstyp, die geänderten Daten und Metadaten zur Quelle. So konsumierst und routest du diese Events.
interface CdcEvent {
op: "c" | "u" | "d" | "r"; // create, update, delete, read (snapshot)
before: Record<string, unknown> | null;
after: Record<string, unknown> | null;
source: {
version: string;
connector: string;
name: string;
ts_ms: number;
db: string;
schema: string;
table: string;
lsn: number;
};
ts_ms: number;
}
type EventHandler = (event: CdcEvent) => Promise<void>;
class CdcEventRouter {
private handlers: Map<string, EventHandler[]> = new Map();
on(table: string, handler: EventHandler): void {
const existing = this.handlers.get(table) ?? [];
existing.push(handler);
this.handlers.set(table, existing);
}
async process(event: CdcEvent): Promise<void> {
const table = event.source.table;
const handlers = this.handlers.get(table) ?? [];
for (const handler of handlers) {
try {
await handler(event);
} catch (error) {
console.error(
`Handler failed for ${table}:`,
error
);
// Dead letter queue for failed events
await this.sendToDeadLetter(event, error as Error);
}
}
}
private async sendToDeadLetter(
event: CdcEvent,
error: Error
): Promise<void> {
console.error(
JSON.stringify({
type: "dead_letter",
table: event.source.table,
operation: event.op,
lsn: event.source.lsn,
error: error.message,
timestamp: new Date().toISOString(),
})
);
}
}
// Wire up handlers for different tables
const router = new CdcEventRouter();
router.on("orders", async (event) => {
if (event.op === "c" || event.op === "u") {
await updateSearchIndex("orders", event.after);
await invalidateCache(`order:${event.after?.id}`);
}
if (event.op === "d") {
await removeFromSearchIndex("orders", event.before?.id);
await invalidateCache(`order:${event.before?.id}`);
}
});
router.on("order_items", async (event) => {
if (event.op === "c") {
await updateAnalytics("new_item", event.after);
}
});Umgang mit Schema-Evolution
Datenbankschemata ändern sich. Spalten werden hinzugefügt, umbenannt oder gelöscht. Deine CDC-Pipeline muss diese Änderungen verarbeiten, ohne nachgelagerte Consumer zu brechen.
interface SchemaVersion {
version: number;
table: string;
columns: Map<string, ColumnDef>;
migrations: SchemaMigration[];
}
interface ColumnDef {
name: string;
type: string;
nullable: boolean;
defaultValue?: unknown;
}
interface SchemaMigration {
fromVersion: number;
toVersion: number;
transform: (record: Record<string, unknown>) => Record<string, unknown>;
}
class SchemaRegistry {
private schemas: Map<string, SchemaVersion[]> = new Map();
register(schema: SchemaVersion): void {
const existing = this.schemas.get(schema.table) ?? [];
existing.push(schema);
existing.sort((a, b) => a.version - b.version);
this.schemas.set(schema.table, existing);
}
evolve(
table: string,
record: Record<string, unknown>,
fromVersion: number,
toVersion: number
): Record<string, unknown> {
const versions = this.schemas.get(table);
if (!versions) return record;
let current = { ...record };
for (const schema of versions) {
for (const migration of schema.migrations) {
if (
migration.fromVersion >= fromVersion &&
migration.toVersion <= toVersion
) {
current = migration.transform(current);
}
}
}
return current;
}
}
// Example: handling a column rename
const registry = new SchemaRegistry();
registry.register({
version: 2,
table: "orders",
columns: new Map([
["id", { name: "id", type: "uuid", nullable: false }],
[
"customer_email",
{ name: "customer_email", type: "text", nullable: false },
],
]),
migrations: [
{
fromVersion: 1,
toVersion: 2,
transform: (record) => {
// Column renamed: email → customer_email
const { email, ...rest } = record;
return { ...rest, customer_email: email };
},
},
],
});Exactly-once-Zustellsemantik
CDC-Events müssen exakt einmal verarbeitet werden. Doppelte Verarbeitung führt zu falschen Zählwerten, doppelten Benachrichtigungen und inkonsistenten Daten.
interface ProcessedEvent {
eventId: string;
lsn: number;
processedAt: Date;
}
class IdempotentProcessor {
private processedEvents: Map<string, ProcessedEvent> = new Map();
private generateEventId(event: CdcEvent): string {
// Unique ID from source position + table + operation
return `${event.source.name}:${event.source.lsn}:${event.source.table}:${event.op}`;
}
async processOnce(
event: CdcEvent,
handler: (event: CdcEvent) => Promise<void>
): Promise<{ processed: boolean; reason?: string }> {
const eventId = this.generateEventId(event);
// Check if already processed
if (this.processedEvents.has(eventId)) {
return {
processed: false,
reason: "duplicate",
};
}
try {
// Process within a transaction that also records
// the event as processed
await handler(event);
this.processedEvents.set(eventId, {
eventId,
lsn: event.source.lsn,
processedAt: new Date(),
});
return { processed: true };
} catch (error) {
// Don't mark as processed — allow retry
throw error;
}
}
getLastProcessedLsn(): number {
let maxLsn = 0;
for (const event of this.processedEvents.values()) {
if (event.lsn > maxLsn) maxLsn = event.lsn;
}
return maxLsn;
}
}Gesundheit der CDC-Pipeline überwachen
Eine CDC-Pipeline, die still zurückfällt, ist schlimmer als gar kein CDC — nachgelagerte Systeme liefern veraltete Daten, ohne es zu merken.
interface CdcPipelineMetrics {
replicationLag: number; // milliseconds
eventsPerSecond: number;
lastProcessedLsn: number;
lastEventTimestamp: Date;
errorCount: number;
deadLetterCount: number;
}
class CdcMonitor {
private metrics: CdcPipelineMetrics = {
replicationLag: 0,
eventsPerSecond: 0,
lastProcessedLsn: 0,
lastEventTimestamp: new Date(),
errorCount: 0,
deadLetterCount: 0,
};
private eventTimestamps: number[] = [];
recordEvent(event: CdcEvent): void {
const now = Date.now();
const eventTime = event.source.ts_ms;
this.metrics.replicationLag = now - eventTime;
this.metrics.lastProcessedLsn = event.source.lsn;
this.metrics.lastEventTimestamp = new Date(eventTime);
// Track throughput over sliding window
this.eventTimestamps.push(now);
const windowStart = now - 60_000;
this.eventTimestamps = this.eventTimestamps.filter(
(t) => t > windowStart
);
this.metrics.eventsPerSecond =
this.eventTimestamps.length / 60;
}
checkHealth(): {
healthy: boolean;
alerts: string[];
} {
const alerts: string[] = [];
if (this.metrics.replicationLag > 30_000) {
alerts.push(
`High replication lag: ${this.metrics.replicationLag}ms`
);
}
const timeSinceLastEvent =
Date.now() - this.metrics.lastEventTimestamp.getTime();
if (timeSinceLastEvent > 300_000) {
alerts.push(
`No events for ${Math.round(timeSinceLastEvent / 1000)}s`
);
}
if (this.metrics.deadLetterCount > 100) {
alerts.push(
`${this.metrics.deadLetterCount} events in dead letter queue`
);
}
return {
healthy: alerts.length === 0,
alerts,
};
}
}Die wichtigsten Erkenntnisse
Change Data Capture verwandelt das Transaktionslog deiner Datenbank in einen Echtzeit-Event-Stream, ohne die Quelle mit Polling-Queries oder Trigger-Overhead zu belasten. Debezium stellt einen produktionsreifen Connector bereit, der das PostgreSQL-WAL oder das MySQL-Binlog liest und strukturierte Events in Kafka veröffentlicht. Route Events nach Tabelle und Operationstyp und behandle Inserts, Updates und Deletes mit eigener Logik für jedes nachgelagerte System. Plane Schema-Evolution von Anfang an ein — registriere Schema-Versionen und definiere Migrations-Transformationen, die ältere Events kompatibel mit neueren Consumern halten. Erzwinge Exactly-once-Verarbeitung durch Idempotenzprüfungen auf Basis der Log-Sequenznummer des Events, damit Wiederholungen keine doppelten Effekte nachgelagert auslösen. Überwache den Replication-Lag kontinuierlich, denn eine CDC-Pipeline, die still zurückfällt, liefert ohne Warnung veraltete Daten — und das ist schlimmer als gar keine Pipeline.


