🎨 Sysprovider Code
Sysprovider LogoWiki
🇪🇸Hosting español para ecommerce

Integración de bases de datos en tiempo real con Apache Kafka y cambio de captura de datos (CDC)

Actualizado el 15 de junio de 2026

¡Por supuesto! Aquí tienes el artículo técnico en Markdown puro, optimizado para SEO y listo para publicar.

La latencia es el enemigo silencioso de la arquitectura de datos moderna. Durante años, la sincronización entre bases de datos se basó en procesos batch (ETL nocturnos) que generaban ventanas de inconsistencia de hasta 24 horas. Hoy, la demanda de bases de datos en tiempo real exige un cambio de paradigma. La combinación de Apache Kafka con Captura de Datos de Cambio (CDC) se ha consolidado como la solución estándar para lograr una sincronización continua, fiable y escalable.

Este artículo explora en profundidad cómo integrar bases de datos en tiempo real utilizando Apache Kafka y CDC. Abordaremos la arquitectura, los patrones de implementación, las mejores prácticas y los desafíos operativos, todo ello con un enfoque práctico y técnico.

¿Qué es el Cambio de Captura de Datos (CDC)?

El cambio de captura de datos (CDC) es un patrón de integración que identifica, captura y entrega los cambios realizados en una base de datos de origen (inserciones, actualizaciones y eliminaciones) a sistemas downstream en tiempo real. En lugar de realizar consultas periódicas (polling), CDC escucha el log de transacciones (WAL en PostgreSQL, Binlog en MySQL, Redo Log en Oracle) de la base de datos.

[INFO] El CDC basado en logs es la única forma de garantizar una captura de cambios sin impacto en el rendimiento de la base de datos de producción. No requiere triggers ni columnas de marca de tiempo adicionales.

Tipos de CDC

Existen dos enfoques principales:

  1. CDC basado en Log (Log-based CDC): Lee directamente el registro de transacciones. Es el más eficiente y el menos intrusivo. Es la opción recomendada para producción.
  2. CDC basado en Consultas (Query-based CDC): Realiza consultas periódicas (ej. SELECT * FROM tabla WHERE modified_at > :ultima_vez). Es más simple de implementar pero añade carga a la base de datos y no captura eliminaciones sin una columna lógica.

Apache Kafka: El Sistema Nervioso del Streaming

Apache Kafka no es solo un bus de mensajería; es una plataforma de streaming distribuida que actúa como el sistema nervioso central de la arquitectura. Proporciona tres capacidades clave para la integración en tiempo real:

  • Alta Throughput: Maneja millones de eventos por segundo.
  • Durabilidad y Replay: Los eventos se almacenan en disco y pueden reproducirse desde cualquier punto en el tiempo.
  • Escalabilidad Horizontal: Se escala añadiendo más brokers y particiones.

La integración CDC + Kafka transforma la base de datos en una fuente de eventos. Cada cambio en una fila se convierte en un mensaje (evento) que se publica en un tópico de Kafka.

Arquitectura de Integración: De la Base de Datos al Stream

El patrón arquitectónico más común para lograr bases de datos en tiempo real con sincronización vía CDC y Kafka es el siguiente:

Los Componentes Clave

  1. Base de Datos Fuente: Puede ser PostgreSQL, MySQL, MongoDB, Oracle, SQL Server, etc.
  2. Conector CDC (Debezium): Un conector de Kafka Connect que se ejecuta como worker. Debezium es el estándar de facto. Se conecta al log de la base de datos y serializa los cambios en formato JSON o Avro.
  3. Apache Kafka (Brokers y Topics): Los eventos de cambio (Change Events) se publican en tópicos específicos, generalmente uno por tabla.
  4. Kafka Connect (Sink Connectors): Consumen los eventos del tópico y los escriben en la base de datos destino.
  5. Base de Datos Destino: Puede ser un Data Warehouse (Snowflake, BigQuery), un motor de búsqueda (Elasticsearch), una caché (Redis) o una réplica.

Flujo de Datos Paso a Paso

  1. Captura: Un UPDATE en la tabla pedidos de PostgreSQL escribe un registro en el WAL.
  2. Lectura: Debezium (Kafka Connect Source) lee ese registro del WAL.
  3. Publicación: Debezium publica un mensaje en el tópico dbserver1.public.pedidos con el payload del cambio (valor antes/después, metadatos de la transacción).
  4. Procesamiento (Opcional): Un Stream Processor (Kafka Streams, ksqlDB) puede filtrar, enriquecer o transformar el evento en tiempo real.
  5. Consumo: Un Sink Connector (ej. confluentinc-connect-elasticsearch) consume el evento y lo indexa en Elasticsearch.
  6. Sincronización: La búsqueda en Elasticsearch refleja el cambio en menos de 1 segundo.

[TIP] Para garantizar el orden de los eventos y la consistencia, utiliza el campo __debezium.source.ts_ms o la clave primaria de la tabla como clave del mensaje de Kafka. Esto asegura que todas las actualizaciones de una misma fila vayan a la misma partición.

Implementación Práctica con Debezium y Kafka Connect

Vamos a ver un ejemplo de configuración para capturar cambios de una base de datos PostgreSQL y enviarlos a un tópico de Kafka.

Requisitos Previos

  • Un clúster de Kafka y Kafka Connect en funcionamiento.
  • PostgreSQL con wal_level = logical y replica_identity = FULL (recomendado).

Configuración del Conector Source (Debezium)

Creamos un conector mediante la API REST de Kafka Connect:

curl -X POST http://localhost:8083/connectors -H "Content-Type: application/json" -d '{
  "name": "postgres-connector-pedidos",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "dbz",
    "database.dbname": "ventas_db",
    "database.server.name": "dbserver1",
    "table.include.list": "public.pedidos",
    "plugin.name": "pgoutput",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "key.converter.schemas.enable": "false",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false"
  }
}'

Explicación de parámetros clave:

  • database.server.name: Prefijo del tópico. El conector creará tópicos como dbserver1.public.pedidos.
  • plugin.name: pgoutput es el plugin lógico nativo de PostgreSQL 10+.
  • transforms.unwrap: Transformación que extrae solo el estado "después" del cambio, simplificando el mensaje.
  • key.converter: Sin esquema, solo JSON plano para simplicidad.

Mensaje de Evento Típico

Cuando se actualiza una fila, el evento en el tópico se ve así (simplificado):

{
  "schema": { ... },
  "payload": {
    "before": {
      "id": 123,
      "estado": "pendiente",
      "total": 150.00
    },
    "after": {
      "id": 123,
      "estado": "pagado",
      "total": 150.00
    },
    "source": {
      "version": "2.4.0.Final",
      "connector": "postgresql",
      "name": "dbserver1",
      "ts_ms": 1712345678000,
      "snapshot": "false",
      "db": "ventas_db",
      "schema": "public",
      "table": "pedidos",
      "txId": 12345,
      "lsn": 12345678
    },
    "op": "u",
    "ts_ms": 1712345678123
  }
}
  • op: "u" indica una actualización.
  • op: "c" para creación.
  • op: "d" para eliminación.
  • op: "r" para lectura inicial (snapshot).

Sincronización y Consistencia en Tiempo Real

La promesa de las bases de datos en tiempo real no se cumple solo con capturar eventos; hay que garantizar que los sistemas destino reflejen el estado correcto.

Estrategias de Sincronización

  1. Sincronización Total (Snapshot + CDC):

    • Al iniciar, el conector toma un snapshot de la tabla completa.
    • Posteriormente, solo captura los cambios incrementales.
    • Ideal para arrancar un nuevo sistema o replicar una base de datos completa.
  2. Sincronización Incremental Pura:

    • El conector comienza desde un punto LSN específico o un timestamp.
    • Útil cuando ya existe un estado base en el destino.
  3. Sincronización Bidireccional (Multi-Master):

    • Requiere manejo de conflictos y es extremadamente compleja.
    • Se suele evitar en favor de arquitecturas "single-writer" por tabla.

Manejo de Esquemas y Evolución

Los esquemas de las bases de datos cambian (ALTER TABLE). Debezium detecta automáticamente cambios en el esquema y actualiza el esquema del evento (si usas Avro con Schema Registry). Si usas JSON plano, tu consumidor debe ser tolerante a campos nuevos.

[WARNING] Si añades una columna NOT NULL sin valor por defecto, Debezium puede fallar al serializar el evento. Planifica los cambios de esquema con cuidado y considera usar un Schema Registry para evolucionar los contratos de forma segura.

Casos de Uso Avanzados y Patrones

1. Materialized Views en Tiempo Real con ksqlDB

En lugar de usar un sink connector, puedes consumir el tópico CDC con ksqlDB para crear una vista materializada que se actualiza continuamente.

CREATE STREAM pedidos_stream (
  id BIGINT,
  estado VARCHAR,
  total DOUBLE
) WITH (
  KAFKA_TOPIC = 'dbserver1.public.pedidos',
  VALUE_FORMAT = 'JSON'
);

CREATE TABLE pedidos_por_estado AS
  SELECT estado, COUNT(*) AS total_pedidos, SUM(total) AS suma_total
  FROM pedidos_stream
  GROUP BY estado
  EMIT CHANGES;

Esta tabla se actualiza en milisegundos con cada nuevo cambio.

2. Cache Invalidation (Redis)

Un sink connector de Redis puede escuchar el tópico CDC y, al detectar un op: "u" o op: "d", invalidar o actualizar la clave en caché.

# Configuración del sink Redis (ejemplo conceptual)
curl -X POST http://localhost:8083/connectors -d '{
  "name": "redis-sink-pedidos",
  "config": {
    "connector.class": "com.github.jcustenborder.kafka.connect.redis.RedisSinkConnector",
    "topics": "dbserver1.public.pedidos",
    "redis.hosts": "redis://localhost:6379",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter"
  }
}'

3. Data Mesh y Eventos de Dominio

CDC convierte la base de datos en un Producto de Datos. Cada tabla publica un flujo de eventos que otros equipos pueden consumir sin acoplamiento. Esto es la base de una arquitectura Data Mesh.

Desafíos Operativos y Buenas Prácticas

Gestión de la Carga Inicial (Snapshot)

Tomar un snapshot de una tabla de 10TB puede saturar la red y la CPU de la base de datos.

  • Solución: Usa snapshot.mode: "when_needed" y programa el snapshot en horas de baja actividad.
  • Solución Avanzada: Realiza un snapshot lógico mediante una réplica de solo lectura.

Tolerancia a Fallos y Exactly-Once Semantics

Kafka Connect, combinado con el offset tracking de Debezium, garantiza Al Menos Una Vez (At-Least-Once). Para lograr Exactly-Once en el sink, necesitas un conector que soporte idempotencia (ej. JDBC Sink Connector con modo upsert y clave primaria).

Monitorización

Es crucial monitorizar:

  • Lag del Consumidor: ¿Cuánto se retrasa el sink respecto al source?
  • Tamaño del Log: ¿Los tópicos están creciendo sin control? Configura políticas de retención adecuadas.
  • Estado del Conector: Kafka Connect expone métricas vía JMX y REST (/connectors/{nombre}/status).

Seguridad

  • Cifrado en Tránsito: Usa TLS para Kafka, Kafka Connect y la base de datos.
  • Autenticación: SASL/SCRAM para Kafka, contraseñas fuertes para Debezium.
  • Enmascaramiento de Datos: Debezium permite usar Single Message Transforms (SMTs) para ocultar campos sensibles (ej. MaskField).

Conclusión

La integración de bases de datos en tiempo real mediante Apache Kafka y cambio de captura de datos (CDC) ya no es una opción experimental; es un pilar fundamental de la arquitectura de datos moderna. Permite romper los silos, reducir la latencia a milisegundos y construir sistemas reactivos y resilientes.

Desde la sincronización de réplicas hasta la alimentación de motores de búsqueda y la creación de pipelines de machine learning en streaming, CDC + Kafka ofrece la columna vertebral perfecta para el streaming de datos empresarial. La clave del éxito reside en una implementación cuidadosa, una monitorización constante y una comprensión profunda de los patrones de consistencia.

[INFO] Si estás empezando, el stack Debezium + Kafka Connect + Confluent Platform (o Red Hat AMQ Streams) es la ruta más rápida y probada. No reinventes la rueda; aprovecha los conectores gestionados.

La era de los datos en reposo ha terminado. Bienvenido a la era del streaming y la sincronización perpetua.

¿Necesitas ayuda?Son dos de nuestros técnicos, Agustín y Mikel, y están disponibles para resolver cualquier problema.

Hablar con ellos ahora
Agustín y Mikel