PostgreSQL Logical Replication in 2026: CDC, Change Streaming, and Microservices Integration
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_sizeand 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 →