pg_recvlogical. Это CLI-инструмент, который поставляется вместе с Postgres и позволяет стримить изменения логической репликации в stdout. Это простой способ понять, какие данные доступны через логическую репликацию, чтобы вы могли строить свои собственные пайплайны «прямо из Postgres».Небольшое предупреждение: это демонстрационный пример. В продакшене лучше использовать более надёжный CDC-инструмент, который умеет обрабатывать повторы, обеспечивать устойчивость, эволюцию схем и т.д. Но
pg_recvlogical отлично подходит, чтобы увидеть сам механизм в действии. Также мы собираемся выполнить команды, которые изменяют базу данных. Убедитесь, что у вас тестовая база, а не продакшен, или используйте Docker, чтобы поднять временный экземпляр Postgres.Теперь к коду:
Чтобы запустить Docker Postgres, готовый к логической репликации, сделайте следующее:
Запустите Postgres в Docker-контейнере с
wal_level=logical:docker run --name pg-test -e POSTGRES_PASSWORD=secret -d -p 5432:5432 postgres:latest -c wal_level=logical
Подключитесь к базе (пароль —
secret) и создайте базу, таблицу и репликационные слоты:docker exec -it pg-test createdb -U postgres mydb
docker exec -it pg-test psql -U postgres mydb
CREATE TABLE orders (
id SERIAL PRIMARY KEY,
product TEXT,
quantity INTEGER
);
SELECT pg_create_logical_replication_slot('my_slot', 'test_decoding');
В другом терминале запустите
pg_recvlogical, чтобы стримить изменения:docker exec -it pg-test pg_recvlogical -U postgres -d mydb --slot=my_slot --start -f -
В терминале, где вы подключены к
mydb, вставляйте и обновляйте данные. Выполняйте каждую команду по очереди и наблюдайте поток изменений от pg_recvlogical:INSERT INTO orders (id, product, quantity) VALUES (1, 'Widget', 3);
UPDATE orders SET quantity = 4 WHERE id = 1;
DELETE FROM orders WHERE id = 1;
Вы увидите изменения в реальном времени:
BEGIN 742
table public.orders: INSERT: id[integer]:1 product[text]:'Widget' quantity[integer]:3
table public.orders: UPDATE: id[integer]:1 product[text]:'Widget' quantity[integer]:4
table public.orders: DELETE: id[integer]:1
COMMIT 742
Это, по сути, то, как работают все CDC-инструменты «под капотом»: они подключаются к логическому репликационному слоту и декодируют изменения WAL в поток событий. Конечно, для продакшен-пайплайнов нужно больше — обработка ошибок, повторы, устойчивость, эволюция схем и т.д., но это и есть основной механизм.
👉 @SQLPortal