Databases

PostgreSQL Logical Replication in 2026: CDC, Change Streaming, and Microservices Integration

Ruslan Ismailov Published 18 min read
P

Introduction: What Is Logical Replication and How Does It Differ from Physical Replication

PostgreSQL supports two fundamentally different types of replication: physical (streaming replication) and logical (logical replication). Physical replication operates at the data block level — it copies WAL segments byte by byte from the primary to the standby. This is a powerful tool for high availability, but it is entirely unsuitable for scenarios where you need to selectively replicate individual tables, deliver changes to external systems, or build event-driven architectures.

Logical Replication, introduced in PostgreSQL 10 and significantly expanded in versions 14–17, works at the level of logical changes: instead of blocks, it operates with rows and INSERT, UPDATE, DELETE, and TRUNCATE operations. This makes it an ideal foundation for Change Data Capture (CDC) — a pattern for capturing data changes in real time that, by 2026, has become the backbone of many distributed microservices-based systems.

Key advantages of Logical Replication over physical replication:

  • Replication of individual tables or subsets of tables
  • Replication across different major versions of PostgreSQL
  • Integration with external systems: Kafka, RabbitMQ, Elasticsearch, and other databases
  • Row and column filtering support (introduced in PostgreSQL 15)
  • Foundation for CDC without OS-level external agents

Architecture: Publications, Subscriptions, and Replication Slots Under the Hood

Logical Replication is built around three key concepts: publication, subscription, and replication slot.

Publication

A publication is a named set of tables on the source side (publisher) whose changes will be published. A publication defines what to replicate. You can specify particular tables, all tables in a schema, or all tables in a database. Starting with PostgreSQL 15, publications support row filtering (WHERE) and column selection.

Subscription

A subscription is an object on the receiver side (subscriber) that describes the connection to the publisher and the list of publications to receive data from. The subscriber applies the received changes to its local tables.

Replication Slot

A replication slot is a mechanism that guarantees WAL segments will not be deleted until the subscriber has processed them. Each subscription creates a replication slot on the publisher. The slot stores a position (LSN — Log Sequence Number) up to which the subscriber has read the WAL. Replication slots are what make reliable CDC possible: even if the consumer is temporarily unavailable, data will not be lost.

Internally, PostgreSQL Logical Replication uses an output plugin to decode the WAL. The standard plugin is pgoutput, included in the core. There are also third-party plugins: wal2json (outputs changes as JSON) and decoderbufs (Protocol Buffers, used by Debezium). In 2026, pgoutput is the recommended choice for most scenarios due to its built-in availability and support for all modern PostgreSQL features.

Setting Up Logical Replication: Configuration and Object Creation

postgresql.conf

To enable Logical Replication, the WAL level must be set to logical:

# postgresql.conf
wal_level = logical
max_replication_slots = 10      # maximum number of replication slots
max_wal_senders = 10            # maximum number of WAL sender processes
wal_keep_size = 1GB             # minimum WAL size to retain
max_logical_replication_workers = 4
max_worker_processes = 16

pg_hba.conf

The subscriber must be allowed to connect to the publisher with replication privileges:

# pg_hba.conf on the 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

Creating a Replication User

-- On the publisher
CREATE ROLE replicator WITH LOGIN REPLICATION PASSWORD 'strong_password';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO replicator;
-- For PostgreSQL 15+ with row/column filtering:
GRANT SELECT ON TABLE orders, products TO replicator;

Creating a Publication

-- Publish all tables (PostgreSQL 10+)
CREATE PUBLICATION my_pub FOR ALL TABLES;

-- Publish specific tables
CREATE PUBLICATION orders_pub FOR TABLE orders, order_items
    WITH (publish = 'insert, update, delete');

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

-- PostgreSQL 15+: column selection
CREATE PUBLICATION orders_pub_cols FOR TABLE orders
    (id, status, total_amount, updated_at);

Creating a Subscription

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

-- Check status
SELECT * FROM pg_stat_subscription;
SELECT * FROM pg_subscription_rel;

Change Data Capture with pgoutput and Debezium

CDC is a pattern in which every change to the database is captured as an event and delivered to interested consumers. In the context of PostgreSQL, CDC is implemented through Logical Replication and replication slots.

pgoutput as the Foundation for CDC

The pgoutput plugin decodes the WAL and transmits changes using the Logical Replication Wire Protocol. It supports message types: Begin, Commit, Relation, Insert, Update, Delete, and Truncate. Starting with PostgreSQL 14, support for StreamStart/StreamStop/StreamCommit was added for streaming large transactions without buffering to disk.

Debezium: A Kafka Connect Connector for PostgreSQL

Debezium is the most popular tool for CDC with PostgreSQL in production environments. It works as a Kafka Connect connector, reads from the replication slot, and publishes events to Kafka topics. Connector configuration:

{
  "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"
  }
}

Each Debezium event has a structure with fields before, after, source (including LSN, timestamp, txId), and op (c — create, u — update, d — delete, r — read/snapshot). This enables building reliable event-driven systems with the ability to replay event history.

Microservices Integration Patterns: Event Delivery Without Kafka

In 2026, not every team wants or is able to maintain a Kafka cluster. PostgreSQL Logical Replication makes it possible to build microservices integration directly, without an additional message broker.

Pattern: Transactional Outbox with CDC

The classic dual-write problem (writing to the DB and sending an event) is solved using the Outbox Pattern. A microservice writes events to an outbox table within the same transaction as the core data. A CDC agent reads from the replication slot and delivers events to subscribers:

-- Outbox table on the 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 for outbox only
CREATE PUBLICATION outbox_pub FOR TABLE outbox
    WITH (publish = 'insert');

-- Example of writing to outbox within a business transaction
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;

Pattern: Direct Logical Replication Between Microservices

Each microservice creates its own replication slot and reads only the tables it needs. This eliminates intermediate systems but requires careful slot management to prevent WAL accumulation. Suitable for systems with moderate load (up to a few thousand transactions per second).

Pattern: Event Aggregator

A single service (Event Aggregator) subscribes to the replication slot, transforms events, and publishes them to an internal bus (Redis Streams, NATS, or even a REST API with webhooks). Other microservices subscribe to the bus rather than directly to PostgreSQL. This reduces the load on the publisher.

Implementing CDC in Go: Reading from a Replication Slot with pgx

The pgx library in Go has built-in support for Logical Replication via the pglogrepl package. Let's walk through a complete implementation of a CDC agent.

Dependencies

// 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
)

Main CDC Agent

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)

    // Create replication slot if it doesn't exist
    _, err = pglogrepl.CreateReplicationSlot(
        ctx, conn, slotName, "pgoutput",
        pglogrepl.CreateReplicationSlotOptions{Temporary: false},
    )
    if err != nil {
        // Slot already exists — this is normal
        log.Printf("Slot already exists or error: %v", err)
    }

    // Get current LSN
    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

    // Start replication
    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:
        }

        // Send Standby Status Update to prevent 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)
        // Publish event here: to 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
}

This CDC agent in Go reads changes directly from the replication slot, decodes rows using PostgreSQL type information, and is ready for integration with any downstream system: Redis Streams, NATS, HTTP webhooks, or a custom queue.

Managing Replication Slots: Monitoring, Cleanup, and Preventing WAL Accumulation

Replication slots are the most dangerous aspect of PostgreSQL Logical Replication. If a subscriber stops reading but the slot is not dropped, PostgreSQL will start accumulating WAL segments, which can lead to disk exhaustion and a complete database outage.

Monitoring Replication Slots

-- Status of all 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;

-- How much WAL has accumulated for each 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;

-- WAL sender process statistics
SELECT * FROM pg_stat_replication;

Alerts and Automatic Cleanup

It is recommended to set up alerts when WAL accumulation exceeds 10 GB for any slot. PostgreSQL 16+ introduced the max_slot_wal_keep_size parameter, which automatically deactivates a slot when the limit is exceeded, preventing disk exhaustion:

# postgresql.conf
# PostgreSQL 13+: limit maximum WAL per slot
max_slot_wal_keep_size = 20GB

# Replication lag monitoring
wal_receiver_status_interval = 10s
wal_receiver_timeout = 60s

Dropping Inactive Slots

-- Drop a specific slot
SELECT pg_drop_replication_slot('stale_slot_name');

-- Find and drop all inactive slots with large lag
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 in Kubernetes: StatefulSets, Persistence, and Failover

Deploying PostgreSQL with Logical Replication in Kubernetes requires special attention to data persistence and replication slot management during failover.

StatefulSet for 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

Operators and Failover

In production Kubernetes environments, it is recommended to use operators: Zalando Postgres Operator or CloudNativePG. CloudNativePG version 1.23+ supports automatic replication slot switchover during primary failover, which is critical for uninterrupted CDC.

For local development with Docker, the configuration is much simpler:

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:

Real-World Case Study: Data Synchronization Between Microservices via CDC

Let's look at the architecture of an e-commerce platform with three microservices: Orders Service, Inventory Service, and Analytics Service. Each has its own PostgreSQL database and must react to changes in the Orders Service.

Data Schema and Publication

-- 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: Reserving Items Based on Events

The Inventory Service runs a CDC agent subscribed only to the outbox table. When it receives an OrderCreated event, it reserves the items in its own database. When it receives OrderCancelled, it releases the reservation. All logic is atomic: if event processing fails, the agent does not acknowledge the LSN, and the event will be re-read on the next startup.

Analytics Service: Updating Aggregates

The Analytics Service subscribes directly to the orders table. On each UPDATE of an order's status, it updates metrics in its analytical database (such as ClickHouse or TimescaleDB). Since analytical queries are idempotent, consumers can safely re-read events after failures.

Idempotency and Exactly-Once Semantics

Distributed systems do not provide exactly-once guarantees without additional mechanisms. To ensure idempotency, each consumer stores a table of processed events in its own database:

-- In each consumer's database
CREATE TABLE processed_events (
    event_id UUID PRIMARY KEY,
    processed_at TIMESTAMPTZ DEFAULT NOW()
);

-- Idempotent event processing
INSERT INTO processed_events (event_id)
VALUES ($1)
ON CONFLICT (event_id) DO NOTHING
RETURNING id;
-- If RETURNING returns a row — process the event
-- If not (conflict) — skip it as already processed

Conclusion

PostgreSQL Logical Replication in 2026 is a mature, production-ready technology for building distributed systems. It covers a wide range of use cases: from simple data synchronization between clusters to complex CDC pipelines with Debezium and Kafka, to direct microservices integration without an additional broker.

Key takeaways for architects and backend developers:

  • Replication slots are the foundation of reliable CDC, but require strict monitoring and WAL lag management
  • pgoutput is the recommended plugin for most scenarios in 2026, including integration with Debezium
  • Outbox Pattern + CDC is the most reliable approach for microservices integration with atomic guarantees
  • Go + pglogrepl is an excellent choice for writing lightweight CDC agents without JVM dependencies
  • In Kubernetes, use operators (CloudNativePG) for correct failover and replication slot management
  • Always configure max_slot_wal_keep_size and WAL lag alerts to prevent emergency situations

A properly built CDC pipeline on PostgreSQL Logical Replication enables truly event-driven systems while keeping PostgreSQL as the single source of truth for your business logic.

Technologies

Tags

Ruslan Ismailov

Senior Web / Backend Developer. Senior web/backend developer with 9 years of experience. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservices, CI/CD. More about me →