Передача данных из транзакционных баз в аналитические — неотъемлемая часть любой современной архитектуры данных. Это также постоянная борьба с хрупкими инструментами, высокими затратами и сложной эксплуатацией. Когда в Snowflake начали строить сервис на базе Postgres, решение этой проблемы естественным образом стало приоритетом номер один.

Ниже — технический разбор того, как устроено data mirroring: как репликация Postgres была переосмыслена с нуля.

Оптимизация репликации Postgres

Postgres — прекрасная операционная база данных, но её возможности по захвату изменений (change data capture, CDC) оставляют желать лучшего. Многие пайплайны в итоге получаются хрупкими, поскольку инструментам репликации приходится справляться со сложным взаимодействием непрерывных данных и изменений схемы, снапшотов и сбоев. Чтобы обеспечить надёжный готовый к использованию опыт для Snowflake Postgres, пришлось переизобретать репликацию Postgres с нуля.

Data mirroring — новая функция Snowflake Postgres в публичном превью для устойчивой репликации данных в Snowflake с низкой стоимостью, низкой задержкой и транзакционной консистентностью. Под капотом изменения пушатся прямо из Postgres в таблицы Apache Iceberg™ транзакционными батчами. Эти батчи автоматически применяются к таблицам в Snowflake — транзакционно и бессерверно.

Простота подхода «транзакционный push в data lake, транзакционное применение в Snowflake, никакой дополнительной инфраструктуры» превращает репликацию из хаотичного процесса с массой сложных сценариев отказа в простой часовой механизм, который будет работать вечно.

Нажимаете кнопку — и таблицы Postgres оказываются в Snowflake.

От pull к push: перенос захвата изменений внутрь Postgres

Change data capture — это процесс захвата изменений из транзакционной базы данных в форме, позволяющей воспроизвести их в другой системе.

В Postgres основной механизм для этого называется логическим декодированием — оно преобразует записи WAL в логические операции insert/update/delete на уровне строк. Эти операции отдаются как поток по сети. С этого момента вся нагрузка перекладывается на клиента.

На практике репликация включает намного больше шагов: backfill, изменения схемы, обработка операций create/add/remove/drop table, новые снапшоты таблиц, перезапуск после сбоя, эффективное слияние изменений, сохранение границ транзакций, правильный размер батчей и так далее. Даже встроенная логическая репликация в Postgres покрывает лишь малую часть этих аспектов.

Одна из проблем pull-подхода на основе логического декодирования в том, что внешняя система, потребляющая изменения, ничего не знает о состоянии Postgres. Например, она не знает, когда произошли изменения схемы, как снапшоты таблиц соотносятся с изменениями — и даже жив ли сам Postgres или просто упала сеть.

Решение этой проблемы довольно простое: пушить изменения из Postgres в data lake, а точнее — в таблицы Iceberg (со сжатым Parquet). Объектные хранилища вроде Amazon S3 масштабируемы, надёжны и уже давно используются для бэкапов Postgres. Это правильное место назначения и для захвата изменений.

Data mirroring использует новое расширение Postgres под названием snowflake_cdc, которое непрерывно пушит батчи изменений в пофайловые change log каждой таблицы и в «meta log» в фоновом режиме (с помощью base workers). Преимущество использования расширения в том, что оно точно знает, что происходит внутри Postgres. Оно может аккуратно координировать изменения схемы и сложные транзакции DML/DDL. Оно может делать снапшоты, одновременно пушя изменения и синхронизируя снапшоты с изменениями.

Figure 1. Snowflake Postgres Data Mirroring

Push-based подход к захвату изменений избавляет от целого класса инфраструктурных проблем и эффективно разделяет производителя и потребителя данных через объектное хранилище.

Распутывание временных линий репликации

При построении системы репликации, подобной data mirroring, важный аспект — временная линия базы данных. Процессы репликации работают с состоянием базы в недавнем прошлом.

Каждая запись в Postgres фактически проходит четыре стадии, каждая из которых представляет непрерывный процесс, работающий на одной временной линии, но в разный момент времени:

  • Write: запись изменяет таблицу и добавляется в WAL (в «настоящем моменте»)
  • Decode: прошлый WAL транслируется в изменение на уровне строк
  • Capture: прошлые изменения на уровне строк захватываются батчами
  • Apply: прошлые батчи изменений слитно применяются к целевым таблицам

Процесс декодирования опирается на специальный механизм Postgres для чтения таблиц каталога в том виде, в каком они были на момент записи («исторический снапшот»). Так бинарные записи WAL можно интерпретировать как логические изменения строк, даже если таблица к моменту декодирования уже была изменена или удалена. В случае data mirroring записи попадают во временные файлы.

Периодически декодер получает сигнал зафиксировать текущий батч и отправляет процессу захвата сообщение о его готовности. Процесс захвата дописывает финализированные файлы в change log Iceberg, записывает запись в meta log и отслеживает реплицированный LSN. Изменения схемы проходят тот же путь write → decode → capture и могут порождать новые change log.

Figure 2. Snowflake Postgres Data Mirroring

Новые записи meta log и change log появляются в таблицах Iceberg. Процесс применения в Snowflake работает как конечный автомат, выполняющий все инструкции из meta log. Когда операция представляет собой батч изменений (самый распространённый случай), все соседние батчи для каждой таблицы обрабатываются вместе.

Такой подход гарантирует, что изменения схемы корректно вписаны в поток изменений, даже если они произошли в рамках транзакции, которая делала и другие записи. Если из-за неожиданного сбоя WAL был утерян, Postgres может автоматически пушить новые снапшоты и указывать Snowflake на их потребление — на практике это происходит очень редко благодаря использованию failover-слотов.

Транзакции как строительный блок распределённых систем

Системы баз данных умеют скрывать огромное количество сложности, связанной со сбоями системы и оборудования, через один простой примитив — транзакции.

Если транзакция не выполняется из-за сбоя на нижнем уровне системы, ничего не происходит — можно попробовать снова. Если транзакция выполнена успешно, можно быть уверенным, что эта работа не будет проделана повторно.

Разочарование от работы с процессами вроде ETL или CDC возникает именно потому, что транзакции внезапно перестают действовать, и приходится самостоятельно разбираться со всеми возможными режимами отказа.

Первым ответом на эту проблему стал сервис Postgres для data lake — управляемая версия open source расширения pg_lake, теперь доступного в общем доступе. Он даёт Postgres уникальную возможность выполнять транзакции сразу над таблицами Postgres и таблицами Iceberg. ETL обычно требует внешних инструментов, а пользователям приходится проектировать идемпотентность и вести тщательный учёт. Теперь можно просто использовать SQL: удалить из таблицы Postgres, вставить в таблицу Iceberg и закоммитить. После этого данные становятся доступны для запросов в Snowflake. Этот подход универсален, но SQL не очень подходит для сквозной репликации высокочастотных обновлений — именно этот слой и добавляет data mirroring.

Под капотом data mirroring полностью использует pg_lake и реализацию Iceberg в Snowflake. Оно берёт батчи данных и изменений схемы из таблиц Postgres и пушит их в несколько change log в Iceberg в рамках одной транзакции на стороне Postgres. Затем Snowflake объединяет несколько батчей за раз в рамках транзакции на своей стороне. Это значит, что все таблицы Snowflake продвигаются вперёд в рамках одной транзакции — точно до границы транзакции Postgres, сохраняя корректность внешних ключей и join'ов.

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

Высокопроизводительное применение изменений и live views

Транзакции дают ещё одну важную возможность — корректность при масштабировании.

Распространённый подход к репликации — превращать каждую операцию в своего рода upsert. Причина в том, что очень трудно решить проблему консистентности между снапшотами таблиц и изменениями, а также изменениями, которые могут быть воспроизведены повторно после сбоя. Однако у upsert-подхода есть несколько проблем:

  1. На стороне получателя могут возникать несогласованные промежуточные состояния, особенно при добавлении таблицы
  2. Вставки становится очень дорого реплицировать, поскольку их нужно сопоставлять с существующими строками в целевой таблице, что крайне затратно при колоночном хранении
  3. Становится очень сложно или невозможно эффективно объединять недавние изменения с целевой таблицей

Когда репликация превращается в транзакционный процесс, управляемый со стороны Postgres, эти ограничения исчезают. Можно построить поток идеальных удалений и вставок, применяемых ровно один раз, без риска того, что вставка уже присутствовала в снапшоте или другом батче изменений.

Это означает следующее:

  1. Репликация нагруженных вставками рабочих нагрузок (обычно самых больших таблиц) становится крайне быстрой и экономичной, поскольку вставки дописываются, а не апсертятся.
  2. Недавние изменения можно эффективно объединять с существующими данными.

Отсюда и появились live views.

Live views — это функция data mirroring, объединяющая ещё не применённые изменения из change log каждой таблицы с данными в целевых таблицах. Важно, что любые фильтры и проекции в запросе можно продвинуть напрямую в слой хранения и сканирование таблиц — как для Parquet-файлов в change log, так и для базовой таблицы. Иными словами: live views работают быстро.

Figure 3. Snowflake Postgres Data Mirroring

Благодаря live views больше не нужно применять изменения очень часто, чтобы добиться низкой задержки. Даже при редком применении задержка live view останется существенно ниже минуты, а запросы всё равно будут выполняться быстро, с лишь немного большими накладными расходами.

Репликация как часовой механизм

Сочетание push-based CDC, тщательно продуманных временных линий, транзакционных границ на обеих сторонах и live views превращает репликацию из хаотичного процесса в швейцарские часы. Нет внешних коннекторов, которые могут отстать. Нет снапшотов, конфликтующих с изменениями. Нет upsert'ов, которые замедляются по мере роста таблиц. Есть расширение Postgres, пушащее батчи в объектное хранилище, и Snowflake, применяющий их — оба транзакционно, оба независимо, оба бесконечно.

Настраиваете один раз. Работает вечно.

Data mirroring для Snowflake Postgres

Исторически операционные и аналитические рабочие нагрузки жили в разных мирах. Эти миры соединялись хрупкими пайплайнами и дополнительными системами, увеличивающими затраты, вместе с целым набором сопутствующих проблем. Основная задача — объединить эти миры на прочном фундаменте с обеих сторон с помощью простых, но качественно инженерных решений.

Со Snowflake Postgres доступен production-grade Postgres, которому можно доверять, и теперь два способа объединить рабочие нагрузки:

  • Data mirroring: публичное превью даёт постоянно работающую, автоматическую репликацию из Postgres в Snowflake. Настройка выполняется один раз — таблицы, включая изменения схемы, остаются синхронизированными непрерывно.
  • Postgres для data lake: доступно в общем доступе, даёт гибкое, управляемое разработчиком перемещение данных между Postgres и Snowflake с использованием открытых форматов, таких как Iceberg. Пишете SQL, и данные перемещаются тогда и так, как нужно.

Начало работы со Snowflake Postgres

Для тех, кто готов начать, доступны следующие ссылки: