ADG Оставить заявку
Блог Данные и аналитика 5 мин чтения

PostgreSQL 10 logical replication как CDC-источник для DWH: избавляемся от ночного ETL-батча

Подключаем логическую репликацию PostgreSQL 10 как источник CDC для аналитического хранилища: стриминговая загрузка вместо ночного батча, близкое к реальному времени.

Контекст момента

PostgreSQL 10 logical replication - применение механизма CDC для стриминговой загрузки изменений в DWH без ночного ETL-батча

В феврале мы рассказывали, как использовали логическую репликацию PostgreSQL 10 для безопасной миграции DWH с минимальным даунтаймом. Там она работала как одноразовый инструмент переключения. Сейчас другой случай: тот же механизм - но уже как постоянный канал подачи данных из OLTP в аналитику. Что-то вроде CDC, только без Debezium и Kafka, прямо из PostgreSQL.

Откуда задача

Клиент - розница, транзакционная база на PostgreSQL 10, аналитическое хранилище отдельно. До недавнего времени DWH кормился ночным ETL-батчем: в 02:00 процесс вычитывал изменения через WHERE updated_at > :last_run, перекладывал в staging, прогонял трансформации, обновлял витрины. К утру аналитики видели вчерашние данные.

Задача стала звучать чётче после того, как бизнес захотел смотреть на продажи не к девяти утра, а с задержкой в несколько минут. Дашборды оперативного управления, контроль кассовых операций - там минуты считаются. ETL-батч раз в сутки эту потребность не закрывает.

Варианты рассматривали разные. Уменьшить интервал батча до пятнадцати минут - самое очевидное, но updated_at-подход плохо ловит DELETE и некоторые UPDATE, где меняется только связанная таблица, а не та, по которой смотрим timestamp. Debezium + Kafka - правильная архитектура для крупного масштаба, но для одного источника и одного приёмника это завезти целый зоопарк. Логическая репликация PostgreSQL 10 напрямую - решили попробовать.

Как это устроено

Идея простая: OLTP-база публикует изменения через PUBLICATION, DWH-слой подписывается на них через SUBSCRIPTION. Только вместо того чтобы реплицировать строки напрямую в итоговые таблицы, мы принимаем их в staging-схему, откуда трансформации уже раскладывают по витринам.

На OLTP-источнике:

-- wal_level = logical должен быть в postgresql.conf
CREATE PUBLICATION dwh_feed FOR TABLE
    orders, order_items, customers, products;

На DWH-приёмнике - staging-схема с таблицами, структурно идентичными источнику:

CREATE SUBSCRIPTION dwh_sub
    CONNECTION 'host=oltp-db port=5432 dbname=retail user=repl_user password=...'
    PUBLICATION dwh_feed;

После этого все INSERT/UPDATE/DELETE из указанных таблиц начинают прилетать в staging на DWH почти немедленно. Задержка в нормальных условиях - секунды. Отставание мониторится через pg_stat_subscription.

Поверх staging - отдельный процесс, который разбирает изменения и мёрджит их в аналитические таблицы с нужными трансформациями. Запускается каждые несколько минут. Это и есть узкое горлышко: не репликация, а именно обработка staged-данных.

Что вылезло на практике

Первое - DDL не реплицируется. Это фундаментальное ограничение логической репликации PostgreSQL 10, мы его знали ещё по февральскому опыту. На практике: если в OLTP добавляют колонку, нужно руками добавить её и в staging, иначе репликация встанет с ошибкой. Договорились с командой OLTP: любой DDL через задачу в трекере с уведомлением за сутки. Хрупко, но работает пока проект небольшой.

Второе - удаления требуют REPLICA IDENTITY. По умолчанию PostgreSQL реплицирует DELETE только с первичным ключом - это работает. Но если на таблице нет PK или нужно что-то посложнее, надо выставлять REPLICA IDENTITY FULL, что увеличивает WAL. У клиента все важные таблицы с PK, поэтому не болело.

Третье - начальная синхронизация. При создании подписки PostgreSQL делает начальный COPY всех данных из publication. На таблицах с историей в несколько лет это занимает время - у нас ушло несколько часов. В это время репликация текущих изменений уже идёт параллельно, так что данные не теряются, но надо закладывать окно.

Четвёртое - staging растёт. Логическая репликация льёт все изменения без разбора. Если витрина пересчитывается не мгновенно, staging накапливается. Нужна чистка после успешной обработки, иначе таблицы разрастаются быстро. Добавили процедуру DELETE FROM staging.orders WHERE processed_at IS NOT NULL AND processed_at < now() - interval '1 day', запускается по расписанию.

Схема потока

flowchart LR
    OLTP["OLTP PostgreSQL 10\n(orders, customers...)"] -->|logical replication| STG["DWH staging\n(идентичная схема)"]
    STG -->|трансформации каждые N минут| DWH["DWH витрины\n(аналитические таблицы)"]
    DWH --> BI["BI-система\nдашборды"]

Ограничения подхода

Честно о том, что нас сдерживает. Источник должен быть PostgreSQL 10+ - если OLTP на девятке или на чём-то другом, этот путь не работает напрямую. Только одна СУБД на обоих концах - DWH-приёмник тоже должен быть PostgreSQL. Если хранилище на ClickHouse - нужен промежуточный слой, логическую репликацию прямо туда не бросить. Мониторинг репликационного слота - слот накапливает WAL пока подписчик не считал изменения. Если DWH лёг или завис, слот растёт, диск на OLTP заполняется. Это надо мониторить плотно, иначе OLTP встанет из-за переполнения.

Для задачи этого клиента - один PostgreSQL-источник, PostgreSQL-приёмник, небольшая схема - подход работает и без Kafka. На DWH/BI-проектах держим его как вариант для аналогичных задач где масштаб не требует распределённой очереди.

Где стоим

Пайплайн запущен около трёх недель. Задержка от записи в OLTP до появления в аналитических витринах - от двух до десяти минут в зависимости от нагрузки на трансформации. Бизнес это устраивает: оперативные дашборды теперь показывают текущий день, а не вчерашний.

Ночной батч пока не убрали совсем: он делает полную сверку и пересчитывает агрегаты за прошлые дни на случай если в staging что-то потеряли или порядок обработки нарушился. Страховка. Возможно избыточная, но пока не уверены что стриминговый путь покрывает все edge-кейсы на этой схеме - смотрим.

Контакт

Нужна такая же инженерная работа?

Опишите задачу и контекст. Ответим в течение рабочего дня, при необходимости подпишем NDA.