Bases de datos

PostgreSQL Logical Replication en 2026: CDC, transmisión de cambios e integración con microservicios

Ruslan Ismailov Publicado 18 min de lectura
P

Introducción: qué es la Logical Replication y en qué se diferencia de la física

PostgreSQL admite dos tipos fundamentalmente distintos de replicación: la física (streaming replication) y la lógica (logical replication). La replicación física opera a nivel de bloques de datos: copia byte a byte los segmentos WAL del primary al standby. Es una herramienta poderosa para garantizar alta disponibilidad, pero completamente inadecuada para escenarios donde se necesita replicar tablas específicas, entregar cambios a sistemas externos o construir arquitecturas Event-Driven.

La Logical Replication, introducida en PostgreSQL 10 y ampliada significativamente en las versiones 14–17, opera a nivel de cambios lógicos: en lugar de bloques, trabaja con filas y operaciones INSERT, UPDATE, DELETE, TRUNCATE. Esto la convierte en la base ideal para el Change Data Capture (CDC) — el patrón de captura de cambios de datos en tiempo real que en 2026 se ha convertido en el fundamento de muchos sistemas distribuidos con arquitectura de microservicios.

Ventajas clave de la Logical Replication frente a la física:

  • Replicación de tablas individuales o subconjuntos de ellas
  • Replicación entre distintas versiones principales de PostgreSQL
  • Integración con sistemas externos: Kafka, RabbitMQ, Elasticsearch, otras bases de datos
  • Posibilidad de filtrar filas y columnas (disponible desde PostgreSQL 15)
  • Base para CDC sin agentes externos a nivel de SO

Arquitectura: publication, subscription y replication slots por dentro

La Logical Replication se articula en torno a tres conceptos clave: publication, subscription y replication slot.

Publication

Una publication es un conjunto de tablas con nombre en el lado del origen (publisher) cuyos cambios serán publicados. La publication define qué se replica. Puedes especificar tablas concretas, todas las tablas de un esquema o todas las tablas de la base de datos. A partir de PostgreSQL 15, las publications admiten filtrado de filas (WHERE) y selección de columnas.

Subscription

Una subscription es un objeto en el lado del receptor (subscriber) que describe la conexión al publisher y la lista de publications de las que recibirá datos. El subscriber aplica los cambios recibidos a sus tablas locales.

Replication Slot

Un replication slot es el mecanismo que garantiza que los segmentos WAL no se eliminen hasta que el suscriptor los haya procesado. Cada subscription crea un replication slot en el publisher. El slot almacena la posición (LSN — Log Sequence Number) hasta la que el suscriptor ha leído el WAL. Son precisamente los replication slots los que hacen posible un CDC fiable: aunque el consumidor esté temporalmente inaccesible, los datos no se perderán.

Internamente, la Logical Replication de PostgreSQL utiliza un output plugin para decodificar el WAL. El plugin estándar es pgoutput, incluido en el núcleo. También existen plugins de terceros: wal2json (genera los cambios en JSON), decoderbufs (Protocol Buffers, utilizado por Debezium). En 2026, pgoutput es la opción recomendada para la mayoría de los escenarios gracias a su integración nativa y soporte de todas las funcionalidades modernas de PostgreSQL.

Configuración de la Logical Replication: parámetros y creación de objetos

postgresql.conf

Para habilitar la Logical Replication es necesario establecer el nivel de WAL en logical:

# postgresql.conf
wal_level = logical
max_replication_slots = 10      # número máximo de replication slots
max_wal_senders = 10            # número máximo de procesos WAL sender
wal_keep_size = 1GB             # volumen mínimo de WAL a conservar
max_logical_replication_workers = 4
max_worker_processes = 16

pg_hba.conf

El suscriptor debe tener permiso para conectarse al publisher con privilegios de replicación:

# pg_hba.conf en el publisher
# TYPE  DATABASE        USER            ADDRESS                 METHOD
host    replication     replicator      10.0.0.0/24             scram-sha-256
host    mydb            replicator      10.0.0.0/24             scram-sha-256

Creación del usuario de replicación

-- En el publisher
CREATE ROLE replicator WITH LOGIN REPLICATION PASSWORD 'strong_password';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO replicator;
-- Para PostgreSQL 15+ con filtrado de filas/columnas:
GRANT SELECT ON TABLE orders, products TO replicator;

Creación de la Publication

-- Publicamos todas las tablas (PostgreSQL 10+)
CREATE PUBLICATION my_pub FOR ALL TABLES;

-- Publicamos tablas específicas
CREATE PUBLICATION orders_pub FOR TABLE orders, order_items
    WITH (publish = 'insert, update, delete');

-- PostgreSQL 15+: filtrado de filas
CREATE PUBLICATION active_orders_pub FOR TABLE orders
    WHERE (status != 'archived')
    WITH (publish = 'insert, update, delete');

-- PostgreSQL 15+: selección de columnas
CREATE PUBLICATION orders_pub_cols FOR TABLE orders
    (id, status, total_amount, updated_at);

Creación de la Subscription

-- En el subscriber
CREATE SUBSCRIPTION my_sub
    CONNECTION 'host=publisher-host port=5432 dbname=mydb user=replicator password=strong_password'
    PUBLICATION orders_pub;

-- Verificar el estado
SELECT * FROM pg_stat_subscription;
SELECT * FROM pg_subscription_rel;

Change Data Capture con pgoutput y Debezium

CDC es un patrón en el que cada cambio en la base de datos se registra como un evento y se entrega a los consumidores interesados. En el contexto de PostgreSQL, el CDC se implementa mediante Logical Replication y replication slots.

pgoutput como base para CDC

El plugin pgoutput decodifica el WAL y transmite los cambios según el protocolo Logical Replication Wire Protocol. Admite mensajes de los tipos: Begin, Commit, Relation, Insert, Update, Delete, Truncate. A partir de PostgreSQL 14, se añadió soporte para StreamStart/StreamStop/StreamCommit, que permite la transmisión en streaming de transacciones grandes sin necesidad de almacenamiento en disco.

Debezium: conector Kafka Connect para PostgreSQL

Debezium es la herramienta más popular para CDC con PostgreSQL en entornos de producción. Funciona como conector de Kafka Connect, lee desde el replication slot y publica los eventos en topics de Kafka. Configuración del conector:

{
  "name": "postgres-cdc-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres-host",
    "database.port": "5432",
    "database.user": "replicator",
    "database.password": "strong_password",
    "database.dbname": "mydb",
    "database.server.name": "mydb",
    "plugin.name": "pgoutput",
    "publication.name": "orders_pub",
    "slot.name": "debezium_slot",
    "table.include.list": "public.orders,public.order_items",
    "heartbeat.interval.ms": "5000",
    "slot.drop.on.stop": "false",
    "snapshot.mode": "initial",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.add.fields": "op,table,lsn,source.ts_ms"
  }
}

Cada evento de Debezium tiene una estructura con los campos before, after, source (que incluye LSN, timestamp, txId) y op (c — create, u — update, d — delete, r — read/snapshot). Esto permite construir sistemas Event-Driven fiables con capacidad de reproducir el historial de eventos.

Patrones de integración con microservicios: entrega de eventos sin Kafka

En 2026, no todos los equipos quieren o pueden mantener un clúster de Kafka. La Logical Replication de PostgreSQL permite construir integraciones entre microservicios de forma directa, sin necesidad de un broker de mensajes adicional.

Patrón: Transactional Outbox con CDC

El clásico problema de la doble escritura (escribir en la BD + enviar el evento) se resuelve mediante el Outbox Pattern. El microservicio escribe los eventos en una tabla outbox dentro de la misma transacción que los datos principales. El agente CDC lee desde el replication slot y entrega los eventos a los suscriptores:

-- Tabla Outbox en el publisher
CREATE TABLE outbox (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    aggregate_type VARCHAR(100) NOT NULL,
    aggregate_id UUID NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    payload JSONB NOT NULL,
    created_at TIMESTAMPTZ DEFAULT NOW()
);

-- Publication solo para outbox
CREATE PUBLICATION outbox_pub FOR TABLE outbox
    WITH (publish = 'insert');

-- Ejemplo de escritura en outbox dentro de una transacción de negocio
BEGIN;
  UPDATE orders SET status = 'shipped' WHERE id = $1;
  INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
  VALUES ('Order', $1, 'OrderShipped', jsonb_build_object(
      'order_id', $1,
      'shipped_at', NOW(),
      'tracking_number', $2
  ));
COMMIT;

Patrón: Logical Replication directa entre microservicios

Cada microservicio crea su propio replication slot y lee únicamente las tablas que necesita. Esto elimina sistemas intermediarios, pero requiere una gestión cuidadosa de los slots para evitar la acumulación de WAL. Es adecuado para sistemas con carga moderada (hasta varios miles de transacciones por segundo).

Patrón: Event Aggregator

Un único servicio (Event Aggregator) se suscribe al replication slot, transforma los eventos y los publica en un bus interno (Redis Streams, NATS, o incluso una REST API con webhook). El resto de los microservicios se suscriben al bus, no directamente a PostgreSQL. Esto reduce la carga sobre el publisher.

Implementación de CDC en Go: lectura desde el replication slot con pgx

La biblioteca pgx en Go tiene soporte nativo para Logical Replication a través del paquete pglogrepl. Veamos una implementación completa de un agente CDC.

Dependencias

// go.mod
module cdc-agent

go 1.22

require (
    github.com/jackc/pgx/v5 v5.6.0
    github.com/jackc/pglogrepl v0.0.0-20240307033717-828fbfe908e9
)

Agente CDC principal

package main

import (
    "context"
    "fmt"
    "log"
    "os"
    "os/signal"
    "syscall"
    "time"

    "github.com/jackc/pglogrepl"
    "github.com/jackc/pgx/v5/pgconn"
    "github.com/jackc/pgx/v5/pgproto3"
    "github.com/jackc/pgx/v5/pgtype"
)

const (
    connStr         = "postgres://replicator:strong_password@localhost:5432/mydb?replication=database"
    slotName        = "cdc_agent_slot"
    publicationName = "orders_pub"
    standbyTimeout  = 10 * time.Second
)

func main() {
    ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
    defer cancel()

    conn, err := pgconn.Connect(ctx, connStr)
    if err != nil {
        log.Fatalf("Failed to connect: %v", err)
    }
    defer conn.Close(ctx)

    // Creamos el replication slot si no existe
    _, err = pglogrepl.CreateReplicationSlot(
        ctx, conn, slotName, "pgoutput",
        pglogrepl.CreateReplicationSlotOptions{Temporary: false},
    )
    if err != nil {
        // El slot ya existe — es normal
        log.Printf("Slot already exists or error: %v", err)
    }

    // Obtenemos el LSN actual
    sysident, err := pglogrepl.IdentifySystem(ctx, conn)
    if err != nil {
        log.Fatalf("IdentifySystem failed: %v", err)
    }
    log.Printf("System ID: %s, Timeline: %d, XLogPos: %s",
        sysident.SystemID, sysident.Timeline, sysident.XLogPos)

    startLSN := sysident.XLogPos

    // Iniciamos la replicación
    err = pglogrepl.StartReplication(ctx, conn, slotName, startLSN,
        pglogrepl.StartReplicationOptions{
            PluginArgs: []string{
                "proto_version '2'",
                fmt.Sprintf("publication_names '%s'", publicationName),
                "messages 'true'",
                "streaming 'true'",
            },
        },
    )
    if err != nil {
        log.Fatalf("StartReplication failed: %v", err)
    }

    relations := map[uint32]*pglogrepl.RelationMessageV2{}
    typeMap := pgtype.NewMap()
    clientXLogPos := startLSN
    nextStandbyMessageDeadline := time.Now().Add(standbyTimeout)

    for {
        select {
        case <-ctx.Done():
            log.Println("Shutting down CDC agent")
            return
        default:
        }

        // Enviamos Standby Status Update para evitar timeout
        if time.Now().After(nextStandbyMessageDeadline) {
            err = pglogrepl.SendStandbyStatusUpdate(ctx, conn,
                pglogrepl.StandbyStatusUpdate{WALWritePosition: clientXLogPos})
            if err != nil {
                log.Printf("SendStandbyStatusUpdate error: %v", err)
            }
            nextStandbyMessageDeadline = time.Now().Add(standbyTimeout)
        }

        ctx2, cancel2 := context.WithDeadline(ctx, nextStandbyMessageDeadline)
        rawMsg, err := conn.ReceiveMessage(ctx2)
        cancel2()

        if err != nil {
            if pgconn.Timeout(err) {
                continue
            }
            log.Printf("ReceiveMessage error: %v", err)
            continue
        }

        if errMsg, ok := rawMsg.(*pgproto3.ErrorResponse); ok {
            log.Printf("PostgreSQL WAL error: %+v", errMsg)
            continue
        }

        msg, ok := rawMsg.(*pgproto3.CopyData)
        if !ok {
            log.Printf("Unexpected message type: %T", rawMsg)
            continue
        }

        switch msg.Data[0] {
        case pglogrepl.PrimaryKeepaliveMessageByteID:
            pkm, err := pglogrepl.ParsePrimaryKeepaliveMessage(msg.Data[1:])
            if err != nil {
                log.Printf("ParsePrimaryKeepaliveMessage error: %v", err)
                continue
            }
            if pkm.ReplyRequested {
                nextStandbyMessageDeadline = time.Now()
            }

        case pglogrepl.XLogDataByteID:
            xld, err := pglogrepl.ParseXLogData(msg.Data[1:])
            if err != nil {
                log.Printf("ParseXLogData error: %v", err)
                continue
            }

            logicalMsg, err := pglogrepl.ParseV2(xld.WALData, false)
            if err != nil {
                log.Printf("ParseV2 error: %v", err)
                continue
            }

            processMessage(logicalMsg, relations, typeMap)

            clientXLogPos = xld.WALStart + pglogrepl.LSN(len(xld.WALData))
        }
    }
}

func processMessage(
    msg pglogrepl.Message,
    relations map[uint32]*pglogrepl.RelationMessageV2,
    typeMap *pgtype.Map,
) {
    switch m := msg.(type) {
    case *pglogrepl.RelationMessageV2:
        relations[m.RelationID] = m
        log.Printf("Relation: %s.%s", m.Namespace, m.RelationName)

    case *pglogrepl.InsertMessageV2:
        rel, ok := relations[m.RelationID]
        if !ok {
            log.Printf("Unknown relation ID: %d", m.RelationID)
            return
        }
        values := decodeRow(rel, m.Tuple, typeMap)
        log.Printf("INSERT into %s: %v", rel.RelationName, values)
        // Aquí publicamos el evento: en Redis, HTTP, NATS, etc.

    case *pglogrepl.UpdateMessageV2:
        rel, ok := relations[m.RelationID]
        if !ok {
            return
        }
        newValues := decodeRow(rel, m.NewTuple, typeMap)
        log.Printf("UPDATE in %s: %v", rel.RelationName, newValues)

    case *pglogrepl.DeleteMessageV2:
        rel, ok := relations[m.RelationID]
        if !ok {
            return
        }
        log.Printf("DELETE from %s (OldTuple: %v)",
            rel.RelationName, m.OldTuple != nil)

    case *pglogrepl.BeginMessage:
        log.Printf("BEGIN xid=%d, lsn=%s", m.Xid, m.FinalLSN)

    case *pglogrepl.CommitMessage:
        log.Printf("COMMIT at %s", m.CommitLSN)
    }
}

func decodeRow(
    rel *pglogrepl.RelationMessageV2,
    tuple *pglogrepl.TupleData,
    typeMap *pgtype.Map,
) map[string]interface{} {
    values := make(map[string]interface{})
    if tuple == nil {
        return values
    }
    for i, col := range tuple.Columns {
        colName := rel.Columns[i].Name
        switch col.DataType {
        case 'n': // NULL
            values[colName] = nil
        case 't': // text
            dt, ok := typeMap.TypeForOID(rel.Columns[i].DataTypeOID)
            if !ok {
                values[colName] = string(col.Data)
            } else {
                var v interface{}
                err := dt.Codec.DecodeDatabaseSQLValue(
                    typeMap, rel.Columns[i].DataTypeOID,
                    pgtype.TextFormatCode, col.Data, &v,
                )
                if err != nil {
                    values[colName] = string(col.Data)
                } else {
                    values[colName] = v
                }
            }
        }
    }
    return values
}

Este agente CDC en Go lee los cambios directamente desde el replication slot, decodifica las filas teniendo en cuenta los tipos de PostgreSQL y está listo para integrarse con cualquier sistema downstream: Redis Streams, NATS, HTTP webhook o una cola propia.

Gestión de replication slots: monitoreo, limpieza y prevención de acumulación de WAL

Los replication slots son el punto más crítico de la Logical Replication en PostgreSQL. Si un suscriptor deja de leer y el slot no se elimina, PostgreSQL comenzará a acumular segmentos WAL, lo que puede llenar el disco y detener toda la base de datos.

Monitoreo de replication slots

-- Estado de todos los replication slots
SELECT
    slot_name,
    plugin,
    slot_type,
    active,
    active_pid,
    restart_lsn,
    confirmed_flush_lsn,
    pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS wal_lag
FROM pg_replication_slots
ORDER BY wal_lag DESC;

-- Cuánto WAL se ha acumulado para cada slot
SELECT
    slot_name,
    pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn) AS bytes_behind,
    pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS human_readable
FROM pg_replication_slots
WHERE active = false;

-- Estadísticas de los procesos WAL sender
SELECT * FROM pg_stat_replication;

Alertas y limpieza automática

Se recomienda configurar alertas cuando la acumulación de WAL supere los 10 GB en cualquier slot. En PostgreSQL 16+ se introdujo el parámetro max_slot_wal_keep_size, que desactiva automáticamente el slot al superar el límite, evitando que el disco se llene:

# postgresql.conf
# PostgreSQL 13+: limitamos el WAL máximo por slot
max_slot_wal_keep_size = 20GB

# Monitoreo del retraso de replicación
wal_receiver_status_interval = 10s
wal_receiver_timeout = 60s

Eliminación de slots inactivos

-- Eliminar un slot específico
SELECT pg_drop_replication_slot('stale_slot_name');

-- Encontrar y eliminar todos los slots inactivos con gran retraso
DO $$
DECLARE
    slot RECORD;
BEGIN
    FOR slot IN
        SELECT slot_name
        FROM pg_replication_slots
        WHERE active = false
          AND pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) > 10 * 1024 * 1024 * 1024
    LOOP
        PERFORM pg_drop_replication_slot(slot.slot_name);
        RAISE NOTICE 'Dropped stale slot: %', slot.slot_name;
    END LOOP;
 END
$$;

Logical Replication en Kubernetes: StatefulSets, persistencia y failover

Desplegar PostgreSQL con Logical Replication en Kubernetes requiere especial atención a la persistencia de datos y la gestión de replication slots durante el failover.

StatefulSet para PostgreSQL Publisher

apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: postgres-publisher
  namespace: data-platform
spec:
  serviceName: postgres-publisher
  replicas: 1
  selector:
    matchLabels:
      app: postgres-publisher
  template:
    metadata:
      labels:
        app: postgres-publisher
    spec:
      containers:
      - name: postgres
        image: postgres:17
        env:
        - name: POSTGRES_DB
          value: mydb
        - name: POSTGRES_USER
          valueFrom:
            secretKeyRef:
              name: postgres-secret
              key: username
        - name: POSTGRES_PASSWORD
          valueFrom:
            secretKeyRef:
              name: postgres-secret
              key: password
        - name: POSTGRES_INITDB_ARGS
          value: "--wal-segsize=64"
        ports:
        - containerPort: 5432
        volumeMounts:
        - name: postgres-data
          mountPath: /var/lib/postgresql/data
        - name: postgres-config
          mountPath: /etc/postgresql/postgresql.conf
          subPath: postgresql.conf
        resources:
          requests:
            memory: "2Gi"
            cpu: "1"
          limits:
            memory: "8Gi"
            cpu: "4"
  volumeClaimTemplates:
  - metadata:
      name: postgres-data
    spec:
      accessModes: ["ReadWriteOnce"]
      storageClassName: fast-ssd
      resources:
        requests:
          storage: 100Gi

Operador y failover

En entornos de producción en Kubernetes se recomienda usar operadores: Zalando Postgres Operator o CloudNativePG. CloudNativePG en la versión 1.23+ admite la transferencia automática de replication slots durante el failover del primary, lo cual es crítico para un CDC continuo.

Para desarrollo local con Docker, la configuración es considerablemente más simple:

version: '3.9'
services:
  postgres-publisher:
    image: postgres:17
    environment:
      POSTGRES_DB: mydb
      POSTGRES_USER: admin
      POSTGRES_PASSWORD: secret
    command: >
      postgres
      -c wal_level=logical
      -c max_replication_slots=10
      -c max_wal_senders=10
    ports:
      - "5432:5432"
    volumes:
      - publisher_data:/var/lib/postgresql/data

  cdc-agent:
    build: ./cdc-agent
    environment:
      PG_CONN: postgres://replicator:password@postgres-publisher:5432/mydb
    depends_on:
      - postgres-publisher

volumes:
  publisher_data:

Caso real: sincronización de datos entre microservicios mediante CDC

Veamos la arquitectura de una plataforma de e-commerce con tres microservicios: Orders Service, Inventory Service y Analytics Service. Cada uno tiene su propia base de datos PostgreSQL y debe reaccionar ante los cambios en el Orders Service.

Esquema de datos y publicación

-- Orders Service DB (publisher)
CREATE TABLE orders (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    customer_id UUID NOT NULL,
    status VARCHAR(50) NOT NULL DEFAULT 'pending',
    total_amount DECIMAL(10, 2) NOT NULL,
    created_at TIMESTAMPTZ DEFAULT NOW(),
    updated_at TIMESTAMPTZ DEFAULT NOW()
);

CREATE TABLE order_items (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    order_id UUID NOT NULL REFERENCES orders(id),
    product_id UUID NOT NULL,
    quantity INT NOT NULL,
    unit_price DECIMAL(10, 2) NOT NULL
);

-- CDC Outbox
CREATE TABLE outbox (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    event_type VARCHAR(100) NOT NULL,
    aggregate_id UUID NOT NULL,
    payload JSONB NOT NULL,
    created_at TIMESTAMPTZ DEFAULT NOW()
);

CREATE PUBLICATION orders_cdc_pub
    FOR TABLE orders, order_items, outbox
    WITH (publish = 'insert, update, delete');

Inventory Service: reserva de productos mediante eventos

El Inventory Service ejecuta un agente CDC suscrito únicamente a la tabla outbox. Al recibir el evento OrderCreated, reserva los productos en su propia BD. Ante OrderCancelled, libera la reserva. Toda la lógica es atómica: si el procesamiento del evento falla, el agente no confirma el LSN y el evento será releído en el siguiente arranque.

Analytics Service: actualización de agregados

El Analytics Service está suscrito directamente a la tabla orders. Con cada UPDATE del estado de un pedido, actualiza las métricas en su base de datos analítica (por ejemplo, ClickHouse o TimescaleDB). Dado que las consultas analíticas son idempotentes, el consumidor puede releer eventos de forma segura ante fallos.

Idempotencia y semántica exactly-once

Los sistemas distribuidos no ofrecen garantías exactly-once sin mecanismos adicionales. Para garantizar la idempotencia, cada consumidor almacena en su BD una tabla de eventos procesados:

-- En la BD de cada consumidor
CREATE TABLE processed_events (
    event_id UUID PRIMARY KEY,
    processed_at TIMESTAMPTZ DEFAULT NOW()
);

-- Procesamiento idempotente de un evento
INSERT INTO processed_events (event_id)
VALUES ($1)
ON CONFLICT (event_id) DO NOTHING
RETURNING id;
-- Si RETURNING devuelve una fila — procesamos el evento
-- Si no (conflicto) — lo omitimos como ya procesado

Conclusión

La Logical Replication de PostgreSQL en 2026 es una tecnología madura y lista para producción, orientada a la construcción de sistemas distribuidos. Abarca un amplio espectro de casos de uso: desde la simple sincronización de datos entre clústeres hasta complejos pipelines de CDC con Debezium y Kafka, pasando por la integración directa de microservicios sin broker adicional.

Conclusiones clave para arquitectos y desarrolladores backend:

  • Replication slots — son la base de la fiabilidad del CDC, pero requieren un monitoreo estricto y gestión del lag de WAL
  • pgoutput — es el plugin recomendado para la mayoría de los escenarios en 2026, incluida la integración con Debezium
  • Outbox Pattern + CDC — es la forma más fiable de integrar microservicios con garantías atómicas
  • Go + pglogrepl — una excelente elección para escribir agentes CDC ligeros sin dependencia de la JVM
  • En Kubernetes, utiliza operadores (CloudNativePG) para gestionar correctamente el failover y los replication slots
  • Configura siempre max_slot_wal_keep_size y alertas sobre el WAL lag para prevenir situaciones de emergencia

Un CDC bien construido sobre PostgreSQL Logical Replication permite crear sistemas verdaderamente orientados a eventos, manteniendo PostgreSQL como la única fuente de verdad para tu lógica de negocio.

Tecnologías

Etiquetas

Ruslan Ismailov

Desarrollador Senior Web / Backend. Desarrollador senior web/backend con 9 años de experiencia. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservicios, CI/CD. Más sobre mí →