Журнал, который пишется ради надёжности, оказывается ещё и идеальным источником изменений: в нём по порядку записано всё, что произошло с данными. Из этого простого факта растут три больших применения, которые на первый взгляд кажутся разными инструментами: репликация (отдать журнал другому узлу целиком), CDC (превратить журнал в поток событий для внешних систем) и PITR (восстановиться на любой момент, доиграв журнал до нужной точки). На живом стенде PostgreSQL 18.4 ниже показано, как из одного и того же механизма получаются все три — с реальным pg_stat_replication, реальным декодированным выводом Go-потребителя и Debezium, и реальной строкой лога point-in-time recovery.
Это третья, финальная статья серии «WAL и его аналоги». Опирается на механизм из первой статьи (что такое write-ahead, durability, crash recovery) и на сравнение реализаций из второй (чем WAL PostgreSQL отличается от redo log/binlog MySQL, journal/oplog MongoDB и других).
В статье
- WAL как источник изменений
- Физическая репликация против логической
- Логическое декодирование: сырой слот через Go
- Тот же WAL читает Debezium
- CDC-угол: что Debezium читает из WAL
- PITR: восстановление до точки во времени
- Гарантии и подводные камни
WAL как источник изменений
В первой статье серии WAL разбирался исключительно как механизм durability: журнал нужен, чтобы пережить SIGKILL и восстановить подтверждённые изменения после падения. Это правда, но не вся правда. WAL — это ещё и последовательный, упорядоченный, полный поток всех изменений, которые произошли с данными, с того самого момента, с которого журнал сохранён. А такой поток можно не только проигрывать при восстановлении одного и того же узла — его можно передать на другой узел (репликация), декодировать в события уровня строк для внешних систем (CDC) или доиграть архивную копию до произвольного момента в прошлом (PITR).
Ключевое наблюдение: все три сценария читают один и тот же WAL, просто разными способами и с разными целями. Физическая репликация стримит байты журнала как есть. Логическая репликация и CDC декодируют те же байты в события INSERT/UPDATE/DELETE. PITR берёт архивную копию того же журнала и доигрывает её поверх старого бэкапа. Ничего из этого не требует отдельного, специально выделенного журнала — в отличие, например, от MySQL, где эту работу выполняет второй, физически независимый binlog (разобрано во врезке второй статьи). Экономия PostgreSQL реальна, но не бесплатна: логическое декодирование, физическая репликация и PITR все конкурируют за один и тот же ресурс — WAL-сегменты, которые сервер обязан хранить, пока хотя бы один из потребителей их не подтвердил. Это возвращается в разделе про подводные камни.
Физическая репликация против логической
Физическая репликация — самый прямой способ использовать WAL как источник: реплика (standby) получает от primary поток WAL байт-в-байт и проигрывает его той же машиной восстановления, что применяется при crash recovery. Standby не выполняет SQL заново и не пересчитывает ничего — она буквально накатывает те же изменения страниц, что применил бы primary сам себе. Отсюда и первое следствие: физическая реплика — это всегда точная копия того же кластера той же версии PostgreSQL, потому что физический формат WAL жёстко привязан к внутреннему устройству страниц конкретной версии.
На стенде (wal/replication, 2 инстанса postgres:18.4 — primary на порту 5441, standby на 5442, поднятый через pg_basebackup -R -X stream) это видно напрямую. После INSERT INTO repl VALUES (42) на primary:
primary repl count: 1
standby repl count: 1
=== pg_stat_replication на primary ===
walreceiver|streaming|async|00:00:00.003349
=== standby в recovery? ===
t
count(*)=1 совпадает на обоих узлах — строка появилась на standby без единого SQL-запроса к нему, только за счёт проигрывания WAL. pg_stat_replication на primary показывает application_name=walreceiver, state=streaming (не catchup, не startup — реплика уже в установившемся режиме потоковой передачи) и sync_state=async. А pg_is_in_recovery() на standby возвращает t — реплика физически находится в постоянном режиме восстановления, тот же режим, в котором primary был бы сразу после своего SIGKILL из первой статьи, только здесь этот режим никогда не заканчивается, пока standby остаётся репликой.
sync_state=async — не техническая деталь, а прямой durability-контракт. Асинхронная репликация (дефолт) означает: primary подтверждает COMMIT клиенту, как только WAL зафиксирован локально, не дожидаясь, пока standby подтвердит получение той же записи. Если primary упадёт до того, как последние несколько транзакций долетели до standby, эти транзакции пропадут при промоушене standby в primary — они были подтверждены клиенту, но реплика их не увидела. Синхронная репликация (synchronous_standby_names настроен, sync_state=sync) закрывает эту дыру ценой латентности: COMMIT на primary не возвращается, пока хотя бы одна синхронная реплика не подтвердила запись — то же самое разделение on/local/off из synchronous_commit, разобранное в первой статье, только теперь применительно не к диску, а к сети между узлами.
Логическая репликация решает принципиально другую задачу: не «дать реплике байт-в-байт то же самое», а «декодировать WAL в поток осмысленных изменений строк — Insert/Update/Delete с конкретными значениями колонок». Для этого WAL должен писаться в режиме wal_level=logical (эта же настройка встречалась в первой статье), а декодирование выполняет специальный плагин вывода (output plugin) — на стенде это pgoutput, встроенный в PostgreSQL с версии 10. Логическая репликация не привязана к идентичной версии и схеме на обеих сторонах: подписчик может быть другой версией PostgreSQL, другой таблицей, или вообще не PostgreSQL — именно логическое декодирование лежит в основе и встроенной publication/subscription-репликации PostgreSQL, и любого внешнего CDC-инструмента вроде Debezium, разобранного ниже.
Логическое декодирование: сырой слот через Go
Логическая репликация в PostgreSQL всегда начинается с репликационного слота — серверного объекта, который резервирует WAL от определённой позиции (LSN) и следит, чтобы сервер не удалил ещё не подтверждённые потребителем сегменты. Для декодирования нужны публикация (CREATE PUBLICATION на нужные таблицы) и подписчик, который стримит слот по протоколу репликации и парсит присланные сообщения. Ниже — как это выглядит на уровне приложения, без встроенной SUBSCRIPTION, напрямую через библиотеку pglogrepl (Go, github.com/jackc/pgx/v5 для соединения + github.com/jackc/pglogrepl для протокола).
Основной цикл потребителя: IDENTIFY_SYSTEM → START_REPLICATION со слотом wal_slot и публикацией wal_pub → бесконечный ReceiveMessage, который различает keepalive-сообщения сервера (нужно вовремя отвечать StandbyStatusUpdate, иначе сервер решит, что клиент отвалился) и собственно XLogData — декодированные байты, которые pglogrepl.Parse превращает в типизированные сообщения протокола pgoutput.
for {
if time.Now().After(nextStandbyMessageDeadline) {
if err := sendStandbyStatusUpdate(ctx, conn, clientXLogPos); err != nil {
return fmt.Errorf("send standby status update: %w", err)
}
nextStandbyMessageDeadline = time.Now().Add(standbyMessageTimeout)
}
ctx2, cancel := context.WithDeadline(ctx, nextStandbyMessageDeadline)
rawMsg, err := conn.ReceiveMessage(ctx2)
cancel()
if err != nil {
if pgconn.Timeout(err) {
continue // не ошибка: пора отправить standby-update на следующей итерации
}
// ...
}
var copyData *pgproto3.CopyData
switch m := rawMsg.(type) {
case *pgproto3.CopyData:
copyData = m
// ...
}
switch copyData.Data[0] {
case pglogrepl.PrimaryKeepaliveMessageByteID:
pkm, err := pglogrepl.ParsePrimaryKeepaliveMessage(copyData.Data[1:])
// ... подвинуть clientXLogPos, ответить если ReplyRequested
case pglogrepl.XLogDataByteID:
xld, err := pglogrepl.ParseXLogData(copyData.Data[1:])
// ...
msg, err := pglogrepl.Parse(xld.WALData)
// ...
handleMessage(msg, relations)
}
}Сама диспетчеризация внутри handleMessage — по типу декодированного сообщения pgoutput. RelationMessage присылается перед первым изменением по таблице в рамках сессии и описывает её структуру (имя, колонки) — без него Insert/Update/Delete бессмысленны, поэтому потребитель обязан кешировать эти описания по RelationID:
func handleMessage(msg pglogrepl.Message, relations map[uint32]*pglogrepl.RelationMessage) {
switch m := msg.(type) {
case *pglogrepl.RelationMessage:
relations[m.RelationID] = m
log.Printf("Relation: %s.%s (columns: %s)", m.Namespace, m.RelationName, columnNames(m.Columns))
case *pglogrepl.InsertMessage:
rel, ok := relations[m.RelationID]
// ...
fmt.Printf("Insert %s.%s: %s\n", rel.Namespace, rel.RelationName, formatTuple(rel, m.Tuple))
case *pglogrepl.UpdateMessage:
rel, ok := relations[m.RelationID]
// ...
fmt.Printf("Update %s.%s: %s\n", rel.Namespace, rel.RelationName, formatTuple(rel, m.NewTuple))
}
}На стенде (postgres:18.4, слот wal_slot, публикация wal_pub, plugin pgoutput, таблица orders(id, amount, status)) после двух INSERT и одного UPDATE реальный вывод потребителя:
2026/07/06 20:56:18 system identify: systemID=7659511586584404011 timeline=1 xlogpos=0/2159D48 dbname=waldemo
2026/07/06 20:56:18 logical replication started: slot=wal_slot publication=wal_pub startLSN=0/2159D48
2026/07/06 20:56:30 Relation: public.orders (columns: id, amount, status)
Insert public.orders: id=1, amount=100, status=new
Insert public.orders: id=2, amount=250, status=paid
Update public.orders: id=1, amount=100, status=shipped
Ровно два Insert и один Update, с уже читаемыми значениями колонок — включая обновлённое значение status=shipped в Update, полученное из NewTuple (декодер pgoutput при REPLICA IDENTITY DEFAULT присылает в Update полную новую строку, но не обязательно полную старую — к этому нюансу вернёмся в разделе про Debezium). Важная деталь по таймингу, которую легко упустить при повторении демо: сразу после START_REPLICATION движку логического декодирования иногда требуется до 15–40 секунд на «прогрев» (поиск consistent point в WAL), особенно сразу после создания слота с нуля. Изменения, внесённые раньше, чем потребитель реально начал получать XLogData, просто не попадут в окно наблюдения — стоит дать потребителю поработать хотя бы 15–20 секунд вхолостую, прежде чем выполнять INSERT/UPDATE, которые предполагается увидеть.
Тот же WAL читает Debezium
Написанный руками Go-потребитель — это низкий уровень: пользователю приходится самому вести Relation-кеш, обрабатывать keepalive и превращать декодированные тюплы в осмысленную структуру. Debezium делает ровно то же самое — читает тот же WAL через тот же протокол логической репликации с тем же pgoutput — но industrial-grade инструментом со схемой, сериализацией, snapshot’ом начального состояния и (в полной топологии) доставкой в Kafka. На стенде используется Debezium Embedded Engine (debezium-embedded, версия 3.6.0.Final, JDK 21) — без Kafka вообще, offset и schema history пишутся в локальные файлы, а декодированные события просто печатаются в stdout.
Properties props = new Properties();
props.setProperty("name", "wal-debezium-embedded");
props.setProperty("connector.class", "io.debezium.connector.postgresql.PostgresConnector");
// offset/schema-history — локальные файлы (без Kafka)
props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore");
props.setProperty("offset.storage.file.filename", "/tmp/dbz-offsets.dat");
props.setProperty("offset.flush.interval.ms", "1000");
props.setProperty("schema.history.internal", "io.debezium.storage.file.history.FileSchemaHistory");
props.setProperty("schema.history.internal.file.filename", "/tmp/dbz-schema-history.dat");
// подключение к PG-стенду
props.setProperty("database.hostname", System.getProperty("database.hostname", "postgres"));
props.setProperty("database.port", System.getProperty("database.port", "5432"));
props.setProperty("database.dbname", "waldemo");
props.setProperty("topic.prefix", "wal");
// logical decoding через pgoutput; отдельный слот от Go-потребителя (wal_slot)
props.setProperty("plugin.name", "pgoutput");
props.setProperty("slot.name", "debezium_slot");
props.setProperty("publication.autocreate.mode", "filtered");
props.setProperty("table.include.list", "public.orders");
DebeziumEngine<ChangeEvent<String, String>> engine = DebeziumEngine.create(Json.class)
.using(props)
.notifying(record -> {
System.out.println("CHANGE-EVENT: " + record.value());
})
.build();Обратите внимание: slot.name=debezium_slot — отдельный от wal_slot Go-потребителя. Оба слота читают один и тот же WAL таблицы orders независимо друг от друга — репликационные слоты не мешают друг другу и не делят позицию, каждый ведёт собственный курсор.
После INSERT orders(amount, status) VALUES (100, 'new') реальное событие op:c (create), payload из живого прогона (schema-блок опущен для краткости):
{"before":null,"after":{"id":1,"amount":{"scale":0,"value":"ZA=="},"status":"new"},"source":{"version":"3.6.0.Final","connector":"postgresql","name":"wal","ts_ms":1783373121441,"snapshot":"false","db":"waldemo","schema":"public","table":"orders","txId":768,"lsn":34116352},"transaction":null,"op":"c","ts_ms":1783373121885}
После UPDATE orders SET status='paid' WHERE id=1 — событие op:u (update), в той же транзакции:
{"before":null,"after":{"id":1,"amount":{"scale":0,"value":"ZA=="},"status":"paid"},"source":{"version":"3.6.0.Final","connector":"postgresql","name":"wal","ts_ms":1783373121441,"snapshot":"false","db":"waldemo","schema":"public","table":"orders","txId":768,"lsn":34116584},"transaction":null,"op":"u","ts_ms":1783373121899}
Оба события пойманы в одном прогоне, в пределах примерно 10 секунд после соответствующих SQL-команд — то есть тот же порядок величины «прогрева» слота, что и у Go-потребителя выше. Два нюанса в этих payload’ах стоит расшифровать явно, потому что оба выглядят как аномалия при первом взгляде:
amountсериализован как{"scale":0,"value":"ZA=="}, а не просто100. Это штатныйVariableScaleDecimal— формат Debezium дляnumericбез фиксированногоscale:value— base64 от байтового представления числа ("ZA=="декодируется в значение100),scale— количество знаков после запятой. Так Debezium избегает потери точности при сериализацииNUMERICпроизвольной точности в JSON.before:nullв событииop:u, хотя логически предыдущее значениеstatusсуществовало. Причина не в Debezium, а в самом PostgreSQL: таблица создана сREPLICA IDENTITY DEFAULT(значение по умолчанию), при котором PostgreSQL пишет в WAL старые значения только столбцов первичного ключа — этого достаточно, чтобы найти строку для примененияUPDATE/DELETEна реплике, но недостаточно, чтобы восстановить полное старое состояние строки. Если нужен полноценныйbeforeсо всеми столбцами, на таблице нужно явно выставитьALTER TABLE orders REPLICA IDENTITY FULL— ценой того, что WAL для каждогоUPDATEувеличится (полная старая строка вместо одного ключа).
Прогрессия здесь важна как тезис раздела: сырой протокол логической репликации PostgreSQL один и тот же, что руками через pglogrepl, что через промышленный Debezium — разница только в том, что Debezium берёт на себя схему, форматирование типов (та же VariableScaleDecimal), snapshot начального состояния таблицы при первом старте и стандартизацию событий под downstream-потребителей.
CDC-угол: что Debezium читает из WAL
Debezium, показанный выше в PostgreSQL-варианте, — не PostgreSQL-специфичный инструмент. Тот же коннекторный движок (Kafka Connect source connector или embedded engine) умеет читать логический журнал у нескольких СУБД, каждый раз выбирая тот журнал, который у конкретной базы играет роль логического, экспортируемого потока изменений — тема, разобранная во второй статье серии для каждой БД по отдельности:
- PostgreSQL — logical replication слот через
pgoutput(илиwal2json,decoderbufs) — ровно то, что использовано на стенде выше. - MySQL —
binlogв ROW-формате, а не redo log InnoDB: redo log принципиально не покидает пределы движка хранения (см. врезку про MySQL во второй статье), поэтому Debezium подключается к серверу как реплика и читает именно binlog. - MongoDB —
oplog(или change streams поверх него в новых версиях) — уже идемпотентный, разрешённый diff-журнал, разобранный в разделе про MongoDB второй статьи.
Показанный здесь embedded-режим — учебный масштаб: одна таблица, один процесс, вывод в stdout. Полная промышленная топология Debezium — это Kafka Connect worker, реальные Kafka-топики на таблицу, schema registry, коннекторы на несколько исходных БД и десятки downstream-потребителей — она разобрана отдельно, со своим стендом и своей глубиной, в статье «CDC с Debezium: PostgreSQL → Kafka». Здесь не дублируется ни конфигурация Kafka Connect, ни топология топиков — только факт, что в основе всего этого пайплайна лежит то же самое логическое декодирование WAL, что показано выше вручную.
PITR: восстановление до точки во времени
Point-in-time recovery использует WAL для задачи, прямо противоположной репликации по духу, но идентичной по механике: не «применить журнал на другом живом узле в реальном времени», а «взять старый статический бэкап и доиграть по нему архивный WAL ровно до нужного момента в прошлом». Рецепт: периодический базовый бэкап (pg_basebackup) плюс непрерывный архив WAL-сегментов (archive_command, копирующий каждый завершённый сегмент в отдельное хранилище) — вместе они дают возможность восстановить состояние базы на любой момент между временем бэкапа и последним заархивированным сегментом.
На стенде (wal/pitr, тот же образ postgres:18.4, что и primary в первой статье серии) сценарий: базовый бэкап, затем INSERT строки id=1, затем INSERT строки id=2, и восстановление с recovery_target_time между ними — то есть после коммита первой строки, но до коммита второй:
recovery_target_time = 2026-07-07 05:22:51.391141+00 (после id=1, до id=2)
текущее состояние primary (ожидаем 2 строки): 2
=== лог restore-инстанса ===
2026-07-07 05:23:12.698 UTC LOG: starting point-in-time recovery to 2026-07-07 05:22:51.391141+00
2026-07-07 05:23:12.977 UTC LOG: restored log file "000000010000000000000004" from archive
2026-07-07 05:23:13.010 UTC LOG: consistent recovery state reached at 0/3000120
2026-07-07 05:23:13.019 UTC LOG: recovery stopping before commit of transaction 756, time 2026-07-07 05:22:53.798231+00
2026-07-07 05:23:13.675 UTC LOG: archive recovery complete
=== на восстановленном инстансе ===
1|2026-07-07 05:22:48.680126+00
count(*) на восстановленном инстансе (ожидаем 1): 1
pg_is_in_recovery() (ожидаем f — промоутнулся): f
Ключевая строка лога, дословно:
recovery stopping before commit of transaction 756, time 2026-07-07 05:22:53.798231+00
PostgreSQL нашёл в архивном WAL коммит транзакции с id=2, но остановился перед её применением, потому что её время коммита позже заданного recovery_target_time. Итог ровно тот, что и предполагался: count(*)=1 на восстановленном инстансе, видна только строка id=1, id=2 отсутствует — притом что она физически присутствует в архивном WAL, просто восстановление сознательно её не доиграло. pg_is_in_recovery() возвращает f — инстанс уже промоутнулся из режима восстановления в обычный, доступный на чтение и запись, режим (по умолчанию recovery_target_action=promote делает это автоматически по достижении цели).
Честная оговорка про сам стенд: archive_command учебный — test ! -f /path/%f && cp %p /path/%f. Если файл сегмента с таким именем уже существует в архиве (например, от предыдущего прогона демо), команда возвращает код выхода 1, PostgreSQL считает попытку архивации проваленной и бесконечно ретраит именно этот сегмент — pg_stat_archiver в этом случае покажет archived_count, застывший на месте, и растущий failed_count. В продакшене за это отвечают специализированные инструменты вроде pgBackRest или wal-g: они сами отслеживают, что уже заархивировано (обычно по контрольной сумме или отдельному манифесту), не зависят от простого test -f по имени файла и умеют управлять ретеншеном архива и параллельным сжатием. Эксплуатационная сторона — как встроить WAL-архивирование, базовые бэкапы и восстановление в постоянно работающий кластер с автоматическим failover — разобрана в статье «Postgres HA: Patroni и репликация»готовится, с 29 сентября; здесь PITR показан как факт механики, не как готовый эксплуатационный рецепт.
Гарантии и подводные камни
Все три сценария выше эксплуатируют одно и то же общее свойство WAL: сервер не имеет права удалить сегмент, пока хотя бы один зарегистрированный потребитель его не подтвердил. Это свойство и даёт репликации, CDC и PITR их гарантии — и оно же источник самой опасной эксплуатационной ошибки логической (и физической) репликации.
Переполнение слота. Репликационный слот — это обещание сервера: «не удалю WAL, пока ты его не заберёшь». Если потребитель слота (Go-процесс, Debezium, физическая реплика) остановился, упал или просто отстаёт быстрее, чем генерируется WAL, — PostgreSQL продолжает честно копить непрочитанные сегменты на диске неограниченно, до заполнения всего доступного места. Это не баг, а прямое следствие контракта слота, но именно поэтому мониторинг отставания слота (pg_replication_slots.confirmed_flush_lsn в сравнении с текущим pg_current_wal_lsn()) — обязательная эксплуатационная практика, а не опциональная. Забытый, никогда не подключавшийся тестовый слот — классическая причина внезапного заполнения диска на проде.
Retention журнала. Тесно связано с переполнением: у обычного (не слотового) WAL-архивирования retention определяется archive_command/внешним хранилищем и политикой очистки старых сегментов там; у логического или физического слота — исключительно скоростью потребителя. Оба режима требуют явной политики: сколько WAL допустимо накопить, прежде чем считать отставание аварией, и что делать, если потребитель не вернётся (удалить слот вручную — единственный выход, если процесс, который его создал, потерян безвозвратно).
Идемпотентность потребителя. Логическая репликация и CDC не гарантируют доставку “ровно один раз” в строгом смысле: при рестарте потребителя после сбоя (не зафиксировавшего последний confirmed_flush_lsn) сервер отдаст события заново с последней подтверждённой позиции — некоторые события, уже обработанные потребителем до сбоя, могут прийти повторно. И Go-потребитель на pglogrepl, и Debezium подтверждают позицию сами (в коде выше — sendStandbyStatusUpdate, у Debezium — периодический offset.flush.interval.ms), и оба оставляют потребителю ответственность самому быть готовым к повтору: либо делать downstream-операции идемпотентными (upsert по первичному ключу вместо голого insert), либо дедуплицировать по LSN/txId на своей стороне. Тот же принцип, что уже встречался во второй статье применительно к MongoDB oplog: журнал изменений хранит достаточно данных для идемпотентного повторного применения, но саму идемпотентность обеспечивает потребитель, а не журнал.
Все три сценария этой статьи — физическая репликация, логическая репликация с CDC, PITR — не отдельные технологии, а три разных способа прочитать один и тот же WAL. Понимание этого фундамента снимает добрую половину вопросов вида «а почему у нас распух диск после отключения консьюмера» или «почему Debezium не видит старые значения строки» — оба ответа лежат в устройстве самого журнала, а не в конкретном инструменте поверх него.
Источники
- PostgreSQL Documentation — Chapter 27. High Availability, Load Balancing, and Replication, Chapter 49. Logical Replication, Chapter 26. Backup and Restore — Continuous Archiving and Point-in-Time Recovery (PITR)
- pglogrepl — github.com/jackc/pglogrepl
- Debezium Documentation — PostgreSQL connector, Embedded Engine
- Стенды:
digital-cookbook/databases/wal(logical,debezium,replication,pitr; PostgreSQL 18.4, Debezium 3.6.0.Final)
Комментарии