Zum Inhalt springen

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.

5 Min. Lesezeit
Datenflussdiagramm, das zeigt, wie Datenbank-Transaktionslogs erfasst und über eine CDC-Pipeline an mehrere nachgelagerte Consumer gestreamt werden

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.

ymlyaml
# ❌ 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
tstypescript
// ✅ 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.

tstypescript
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.

tstypescript
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.

tstypescript
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.

tstypescript
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.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX