Saltar al contenido

Implementación de pipelines de Change Data Capture en bases de datos

Construye pipelines CDC que transmiten cambios en tiempo real: captura basada en logs con Debezium, formato de eventos, esquemas y entrega exactly-once.

5 min de lectura
Diagrama de flujo de datos que muestra cómo los logs de transacciones de la base de datos se capturan y se transmiten a múltiples consumidores a través de un pipeline CDC

Hacer polling a una base de datos cada pocos segundos para detectar cambios desperdicia recursos e introduce latencia. Change Data Capture (CDC) lee el propio log de transacciones de la base de datos—el write-ahead log (WAL) en PostgreSQL, el binlog en MySQL—y transmite cada inserción, actualización y eliminación como un evento. Esto te da replicación de datos en tiempo real sin afectar el rendimiento de la base de datos origen.

CDC es la columna vertebral de las arquitecturas de datos modernas: impulsa la actualización de índices de búsqueda, la invalidación de cachés, los pipelines de analítica y la sincronización de datos entre servicios sin acoplar productores con consumidores.

CDC basado en logs con Debezium

Debezium lee el log de transacciones de la base de datos y publica eventos de cambio en Kafka. La base de datos origen no necesita ninguna modificación—sin triggers, sin consultas de polling, sin publicación de eventos a nivel de aplicación.

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",
  },
};

El plugin pgoutput lee el stream de replicación lógica de PostgreSQL. La transformación ExtractNewRecordState aplana el formato de envoltura de Debezium convirtiéndolo en el estado posterior de cada fila, lo que simplifica el procesamiento posterior.

Procesamiento de eventos de cambio

Cada evento CDC contiene el tipo de operación, los datos modificados y metadatos sobre el origen. Así es como se consumen y enrutan estos eventos.

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);
  }
});

Gestión de la evolución del esquema

Los esquemas de bases de datos cambian. Las columnas se añaden, se renombran o se eliminan. Tu pipeline CDC debe manejar estos cambios sin romper a los consumidores.

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 };
      },
    },
  ],
});

Semántica de entrega exactly-once

Los eventos CDC deben procesarse exactamente una vez. El procesamiento duplicado genera conteos incorrectos, notificaciones duplicadas e inconsistencia de datos.

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;
  }
}

Monitorización de la salud del pipeline CDC

Un pipeline CDC que se queda atrás silenciosamente es peor que no tener CDC—los sistemas downstream sirven datos obsoletos sin saberlo.

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,
    };
  }
}

Conclusiones clave

Change Data Capture convierte el log de transacciones de tu base de datos en un stream de eventos en tiempo real sin sobrecargar el origen con consultas de polling ni triggers. Debezium proporciona un conector listo para producción que lee el WAL de PostgreSQL o el binlog de MySQL y publica eventos estructurados en Kafka. Enruta los eventos por tabla y tipo de operación, manejando inserciones, actualizaciones y eliminaciones con lógica distinta para cada sistema downstream. Planifica la evolución del esquema desde el principio—registra versiones de esquema y define transformaciones de migración que mantengan los eventos antiguos compatibles con los consumidores nuevos. Garantiza el procesamiento exactly-once mediante verificaciones de idempotencia basadas en el número de secuencia del log del evento, evitando efectos duplicados en los reintentos. Monitoriza el lag de replicación continuamente, porque un pipeline CDC que se queda atrás silenciosamente entrega datos obsoletos sin avisar, lo cual es peor que no tener ningún pipeline.

Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX