Базы данных

PostgreSQL Logical Replication в 2026 году: CDC, потоковая передача изменений и интеграция с микросервисами

Ruslan Ismailov Опубликовано 18 мин чтения
P

Введение: что такое Logical Replication и чем она отличается от физической

PostgreSQL поддерживает два принципиально разных вида репликации: физическую (streaming replication) и логическую (logical replication). Физическая репликация работает на уровне блоков данных — она побайтово копирует WAL-сегменты с primary на standby. Это мощный инструмент для обеспечения высокой доступности, но он совершенно непригоден для сценариев, где нужно выборочно реплицировать отдельные таблицы, доставлять изменения во внешние системы или строить Event-Driven архитектуры.

Logical Replication, появившаяся в PostgreSQL 10 и существенно расширенная в версиях 14–17, работает на уровне логических изменений: вместо блоков она оперирует строками и операциями INSERT, UPDATE, DELETE, TRUNCATE. Это делает её идеальной основой для Change Data Capture (CDC) — паттерна захвата изменений данных в реальном времени, который в 2026 году стал фундаментом многих распределённых систем на микросервисной архитектуре.

Ключевые преимущества Logical Replication перед физической:

  • Репликация отдельных таблиц или их подмножеств
  • Репликация между разными мажорными версиями PostgreSQL
  • Интеграция с внешними системами: Kafka, RabbitMQ, Elasticsearch, другими СУБД
  • Возможность фильтрации строк и столбцов (появилась в PostgreSQL 15)
  • Основа для CDC без внешних агентов на уровне ОС

Архитектура: publication, subscription, replication slots изнутри

Logical Replication строится вокруг трёх ключевых концепций: publication, subscription и replication slot.

Publication

Publication — это именованный набор таблиц на стороне источника (publisher), изменения в которых будут публиковаться. Publication определяет что реплицировать. Вы можете указать конкретные таблицы, все таблицы схемы или все таблицы базы данных. Начиная с PostgreSQL 15, publication поддерживает фильтрацию строк (WHERE) и выбор столбцов.

Subscription

Subscription — это объект на стороне получателя (subscriber), который описывает подключение к publisher и список publication для получения данных. Subscriber применяет полученные изменения к локальным таблицам.

Replication Slot

Replication slot — это механизм, гарантирующий, что WAL-сегменты не будут удалены до тех пор, пока подписчик их не обработал. Каждая subscription создаёт replication slot на publisher. Слот хранит позицию (LSN — Log Sequence Number), до которой подписчик прочитал WAL. Именно replication slots делают возможным надёжный CDC: даже если потребитель временно недоступен, данные не будут потеряны.

Внутри PostgreSQL Logical Replication использует output plugin для декодирования WAL. Стандартный плагин — pgoutput, включённый в ядро. Также существуют сторонние плагины: wal2json (выводит изменения в JSON), decoderbufs (Protocol Buffers, используется Debezium). В 2026 году pgoutput является рекомендуемым выбором для большинства сценариев благодаря встроенности и поддержке всех современных возможностей PostgreSQL.

Настройка Logical Replication: конфигурация и создание объектов

postgresql.conf

Для включения Logical Replication необходимо установить уровень WAL в logical:

# postgresql.conf
wal_level = logical
max_replication_slots = 10      # максимальное число replication slots
max_wal_senders = 10            # максимальное число WAL sender процессов
wal_keep_size = 1GB             # минимальный объём WAL для хранения
max_logical_replication_workers = 4
max_worker_processes = 16

pg_hba.conf

Подписчик должен иметь право подключаться к publisher с репликационными привилегиями:

# pg_hba.conf на 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

Создание пользователя репликации

-- На publisher
CREATE ROLE replicator WITH LOGIN REPLICATION PASSWORD 'strong_password';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO replicator;
-- Для PostgreSQL 15+ с фильтрацией строк/столбцов:
GRANT SELECT ON TABLE orders, products TO replicator;

Создание Publication

-- Публикуем все таблицы (PostgreSQL 10+)
CREATE PUBLICATION my_pub FOR ALL TABLES;

-- Публикуем конкретные таблицы
CREATE PUBLICATION orders_pub FOR TABLE orders, order_items
    WITH (publish = 'insert, update, delete');

-- PostgreSQL 15+: фильтрация строк
CREATE PUBLICATION active_orders_pub FOR TABLE orders
    WHERE (status != 'archived')
    WITH (publish = 'insert, update, delete');

-- PostgreSQL 15+: выбор столбцов
CREATE PUBLICATION orders_pub_cols FOR TABLE orders
    (id, status, total_amount, updated_at);

Создание Subscription

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

-- Проверить состояние
SELECT * FROM pg_stat_subscription;
SELECT * FROM pg_subscription_rel;

Change Data Capture с pgoutput и Debezium

CDC — это паттерн, при котором каждое изменение в базе данных фиксируется как событие и доставляется заинтересованным потребителям. В контексте PostgreSQL CDC реализуется через Logical Replication и replication slots.

pgoutput как основа CDC

Плагин pgoutput декодирует WAL и передаёт изменения в протоколе Logical Replication Wire Protocol. Он поддерживает сообщения типов: Begin, Commit, Relation, Insert, Update, Delete, Truncate. Начиная с PostgreSQL 14, добавлена поддержка StreamStart/StreamStop/StreamCommit для потоковой передачи больших транзакций без буферизации на диске.

Debezium: Kafka Connect коннектор для PostgreSQL

Debezium — самый популярный инструмент для CDC с PostgreSQL в production-средах. Он работает как Kafka Connect коннектор, читает из replication slot и публикует события в Kafka-топики. Конфигурация коннектора:

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

Каждое событие Debezium имеет структуру с полями before, after, source (включая LSN, timestamp, txId) и op (c — create, u — update, d — delete, r — read/snapshot). Это позволяет строить надёжные Event-Driven системы с возможностью воспроизведения истории событий.

Паттерны интеграции с микросервисами: доставка событий без Kafka

В 2026 году далеко не каждая команда хочет или может поддерживать Kafka-кластер. PostgreSQL Logical Replication позволяет строить интеграцию микросервисов напрямую, без дополнительного брокера сообщений.

Паттерн: Transactional Outbox с CDC

Классическая проблема двойной записи (запись в БД + отправка события) решается через Outbox Pattern. Микросервис пишет события в таблицу outbox в той же транзакции, что и основные данные. CDC-агент читает из replication slot и доставляет события подписчикам:

-- Таблица Outbox на 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 только для outbox
CREATE PUBLICATION outbox_pub FOR TABLE outbox
    WITH (publish = 'insert');

-- Пример записи в outbox внутри бизнес-транзакции
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;

Паттерн: Direct Logical Replication между микросервисами

Каждый микросервис создаёт собственный replication slot и читает только нужные ему таблицы. Это исключает промежуточные системы, но требует тщательного управления слотами для предотвращения накопления WAL. Подходит для систем с умеренной нагрузкой (до нескольких тысяч транзакций в секунду).

Паттерн: Event Aggregator

Один сервис (Event Aggregator) подписывается на replication slot, трансформирует события и публикует их во внутреннюю шину (Redis Streams, NATS, или даже REST API с webhook). Остальные микросервисы подписываются на шину, а не напрямую на PostgreSQL. Это снижает нагрузку на publisher.

Реализация CDC на Go: чтение из replication slot с pgx

Библиотека pgx в Go имеет встроенную поддержку Logical Replication через пакет pglogrepl. Рассмотрим полноценную реализацию CDC-агента.

Зависимости

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

Основной CDC-агент

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)

    // Создаём replication slot если не существует
    _, err = pglogrepl.CreateReplicationSlot(
        ctx, conn, slotName, "pgoutput",
        pglogrepl.CreateReplicationSlotOptions{Temporary: false},
    )
    if err != nil {
        // Слот уже существует — это нормально
        log.Printf("Slot already exists or error: %v", err)
    }

    // Получаем текущий 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

    // Запускаем репликацию
    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:
        }

        // Отправляем Standby Status Update для предотвращения 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)
        // Здесь публикуем событие: в Redis, HTTP, NATS и т.д.

    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
}

Этот CDC-агент на Go читает изменения из replication slot напрямую, декодирует строки с учётом типов PostgreSQL и готов к интеграции с любой downstream-системой: Redis Streams, NATS, HTTP webhook, или собственной очередью.

Управление replication slots: мониторинг, очистка и предотвращение WAL-накопления

Replication slots — самое опасное место в PostgreSQL Logical Replication. Если подписчик перестал читать, а слот не удалён, PostgreSQL начнёт накапливать WAL-сегменты, что может привести к заполнению диска и остановке всей базы.

Мониторинг replication slots

-- Состояние всех 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;

-- Сколько WAL накоплено для каждого слота
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 процессов
SELECT * FROM pg_stat_replication;

Алерты и автоматическая очистка

Рекомендуется настроить алерты на накопление WAL более 10 ГБ для любого слота. В PostgreSQL 16+ появился параметр max_slot_wal_keep_size, который автоматически деактивирует слот при превышении лимита, предотвращая заполнение диска:

# postgresql.conf
# PostgreSQL 13+: ограничиваем максимальный WAL на слот
max_slot_wal_keep_size = 20GB

# Мониторинг задержки репликации
wal_receiver_status_interval = 10s
wal_receiver_timeout = 60s

Удаление неактивных слотов

-- Удалить конкретный слот
SELECT pg_drop_replication_slot('stale_slot_name');

-- Найти и удалить все неактивные слоты с большим лагом
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 в Kubernetes: StatefulSets, персистентность и failover

Развёртывание PostgreSQL с Logical Replication в Kubernetes требует особого внимания к персистентности данных и управлению replication slots при failover.

StatefulSet для 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

Оператор и failover

В production-среде в Kubernetes рекомендуется использовать операторы: Zalando Postgres Operator или CloudNativePG. CloudNativePG в версии 1.23+ поддерживает автоматическое переключение replication slots при failover primary, что критично для непрерывного CDC.

При использовании Docker для локальной разработки конфигурация значительно проще:

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:

Реальный кейс: синхронизация данных между микросервисами через CDC

Рассмотрим архитектуру e-commerce платформы с тремя микросервисами: Orders Service, Inventory Service и Analytics Service. Каждый из них имеет собственную PostgreSQL базу данных и должен реагировать на изменения в Orders Service.

Схема данных и публикация

-- 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: резервирование товаров по событиям

Inventory Service запускает CDC-агент, подписанный только на таблицу outbox. При получении события OrderCreated он резервирует товары в своей БД. При OrderCancelled — освобождает резерв. Вся логика атомарна: если обработка события падает, агент не подтверждает LSN, и событие будет перечитано при следующем запуске.

Analytics Service: обновление агрегатов

Analytics Service подписан на таблицу orders напрямую. При каждом UPDATE статуса заказа он обновляет метрики в своей аналитической БД (например, ClickHouse или TimescaleDB). Поскольку аналитические запросы идемпотентны, потребитель может безопасно перечитывать события при сбоях.

Идемпотентность и exactly-once семантика

Distributed systems не дают гарантий exactly-once без дополнительных механизмов. Для обеспечения идемпотентности каждый потребитель хранит в своей БД таблицу обработанных событий:

-- В БД каждого потребителя
CREATE TABLE processed_events (
    event_id UUID PRIMARY KEY,
    processed_at TIMESTAMPTZ DEFAULT NOW()
);

-- Идемпотентная обработка события
INSERT INTO processed_events (event_id)
VALUES ($1)
ON CONFLICT (event_id) DO NOTHING
RETURNING id;
-- Если RETURNING вернул строку — обрабатываем событие
-- Если нет (конфликт) — пропускаем как уже обработанное

Заключение

PostgreSQL Logical Replication в 2026 году — это зрелая, production-ready технология для построения распределённых систем. Она покрывает широкий спектр задач: от простой синхронизации данных между кластерами до сложных CDC-пайплайнов с Debezium и Kafka, и до прямой интеграции микросервисов без дополнительного брокера.

Ключевые выводы для архитекторов и backend-разработчиков:

  • Replication slots — основа надёжности CDC, но требуют строгого мониторинга и управления лагом WAL
  • pgoutput — рекомендуемый плагин для большинства сценариев в 2026 году, включая интеграцию с Debezium
  • Outbox Pattern + CDC — наиболее надёжный способ интеграции микросервисов с атомарными гарантиями
  • Go + pglogrepl — отличный выбор для написания легковесных CDC-агентов без зависимости от JVM
  • В Kubernetes используйте операторы (CloudNativePG) для корректного управления failover и replication slots
  • Всегда настраивайте max_slot_wal_keep_size и алерты на WAL lag для предотвращения аварийных ситуаций

Правильно выстроенный CDC на PostgreSQL Logical Replication позволяет строить по-настоящему событийно-ориентированные системы, сохраняя PostgreSQL как единственный источник истины для вашей бизнес-логики.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →