PRODСуверенная европейская платформа BaaSОткрыть панель управления →

Инженерное дело · 10 минута чтения

PostgreSQL CDC с NATS: архитектура реального времени

Affane Daylami · Fondateur · 24 июля 2026 г.

Вернуться в блог

Клиент Aurabase, подключенный через WebSocket, видит, что строка, вставленная в базу, появляется через несколько миллисекунд после COMMIT — без опроса и без веб-перехватчика для настройки. Механизм состоит из трех этапов. Плагин логической репликации декодирует Postgres WAL в JSON, право на его чтение имеет только один процесс, а NATS ретранслирует каждое событие на все серверы WebSocket в кластере. Вот как это реализовано в реальном коде сервиса aura-realtime.

Этот текст на английском языке был создан автоматически на основе французского оригинала и еще не проверялся.
Эта страница была переведена автоматически. Английская версия является авторитетной.

Ничто здесь не является идеализированным образцом. Каждый экстракт поставляется в исходном виде из ящика aura-realtimeна дату публикации. Здесь есть четкое различие: NATS JetStream не является тем, что передает события CDC в этом конвейере. Ниже мы подробно расскажем, что здесь на самом деле делает JetStream и что вместо него переносит трафик CDC. Проверенная ловушка Kubernetes, способная заставить весь механизм молчать без каких-либо ошибок, описана ниже.

Самое необходимое
  • Postgres CDC Aurabase читает WAL через wal2json и pg_logical_slot_get_changes — а не двоичный протокол pgoutputи не мост Debezium/Kafka Connect.
  • Логический слот aura_cdc_slot допускает только один считыватель: aura-realtime-cdc-worker выбирается посредством аренды Kubernetes (coordination.k8s.io/v1), с режимом «всегда лидер», исключая K8s.
  • Разветвление на реплики ws-front осуществляется через ядро ​​NATS (простая публикация/подписка на aura.realtime.>), а не через постоянный поток JetStream. JetStream в этой же службе обслуживает исключительно KV межэкземплярного присутствия.
  • Каждое событие перепроверяется подписчиком в RLS непосредственно перед передачей посредством реального запроса SET LOCAL ROLE на соответствующей линии, а не кэшированной аппроксимации.
  • Проверенная ошибка в истории репозитория: без плагина wal2json в образе Postgres Kubernetes создание слота завершается сбоем, и CDC остается молча неактивным.
#
Выбор двигателя

wal2json, а не pgoutput, а не Debezium

Большинство конвейеров CDC Postgres проходят через pgoutput, протокол репликации двоичной логики, а затем через соединитель, такой как Debezium, который преобразует его в Kafka. Aurabase пропускает этот шаг: сервис aura-realtime напрямую декодирует WAL в логическую репликацию с плагином wal2json, который создает пригодный для использования JSON без промежуточного этапа трансляции.

cdc/postgres.rs (реальные запросы)sql
-- Создается один раз, если отсутствует, воссоздается, если срок действия WAL истек.
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Повторный опрос (задержка 100 мс, если пакет возвращается пустым)
SELECT data FROM pg_logical_slot_get_changes(
  'aura_cdc_slot', NULL, 1000,
  'include-xids', '0', 'format-version', '2',
  'include-schemas', '1', 'include-types', '0');

Три обязательных условия Postgres, проверенные в коде: роль приложения должна содержать REPLICATION, wal_level = logical должна быть активной и max_slot_wal_keep_size должна ограничивать сохранение WAL. Без этого последнего ограничения медленный потребитель приводит к бесконечному росту тома диска. Служба также отслеживает wal_status слота: если он меняется на lost (WAL очищается за пределами этого предела), слот автоматически создается заново. Предполагается и регистрируется потеря событий, не использованных в это время.

Пакет из 1000 изменений на опрос с динамическим фильтром таблицы, обновляемым каждые 3 секунды. Только таблицы, в которых проект явно включил режим реального времени, вводят add-tables — это позволяет избежать декодирования WAL таблиц, которые не представляют интереса для подписчиков.

#
Один продюсер

cdc-worker: выбирается по договору аренды Kubernetes, не дублируется.

Слот логической репликации Postgres допускает одновременное использование только одного активного диска. Параллельное выполнение нескольких потребителей одного и того же aura_cdc_slot нарушит порядок изменений, а не просто продублирует их. Aurabase разделила aura-realtime на два отдельных двоичных файла, чтобы устранить это ограничение, не жертвуя при этом горизонтальным масштабом WebSocket.

Cargo.tomltoml
# cdc-worker: избранный лидер, запускает CDC + публикует на NATS. Нет WS-сервера.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N реплик (HPA), использует NATS, фильтр RLS, обслуживает WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker получает право на чтение слота, только удерживая арендуемый объект Kubernetes (coordination.k8s.io/v1). Он пытается создать его или украсть, если срок его действия истек, а затем обновляет его через треть срока его жизни. Когда аренда потеряна (не удалось продлить срок аренды, сменился владелец), CancellationToken прерывает незавершенную работу, и процесс возвращается в цикл приобретения. За пределами Kubernetes — например, во время локального тестирования или некластерной разработки — клиент K8s не может подключиться, и сервис переключается в режим «всегда лидирующий». Полезно в дев, неправильно, если забыть в проде с несколькими репликами.

cdc/leader.rs (настоящие подписи)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // обновляет все Lease_duration_secs / 3
  // за исключением K8: работает (токен) только один раз, токен никогда не отменяется
}
#
Настоящее развлечение

Ядро NATS для CDC, JetStream для присутствия

Это наиболее часто неправильно понимаемый момент в таком конвейере: NATS и JetStream — это две разные вещи, и события CDC не проходят через постоянный поток JetStream. cdc-worker публикует каждое событие с помощью Client::publish_with_headers — основного API публикации/подписки NATS, который доставляет не более одного раза без сохранения или воспроизведения, отдельно от JetStream.

Сердце НАТС (паб/саб)Распределение событий CDC на реплики ws-frontМаксимум один раз, без настойчивости или повторения
НАТС ДжетСтрим (КВ)Межэкземплярное присутствие (aura_presence)Общее состояние 60-х, а не поток CDC для воспроизведения
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker публикует (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front прослушивает все поддерево (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Каждая реплика ws-front подписывается на поддерево aura.realtime.>и декодирует тему в исходное имя канала. Затем он повторно публикует событие в своем локальном tokio::sync::broadcastдля подключенных к нему WebSockets. Заголовок Aura-Origin содержит идентификатор отправляющего экземпляра: каждая реплика игнорирует сообщения, которые она сама опубликовала, что позволяет избежать дублирования без центральной координации.

Этот выбор имеет прямое следствие: ядро ​​NATS осуществляет доставку в лучшем случае один раз (не более одного раза) без длительной очереди. Если реплика ws-front на короткое время отключается от NATS при прохождении события, она не догоняет его. История каждого канала (Channel.history, буфер в памяти, ограниченный HISTORY_SIZE) хранится локально в каждой реплике, а не на уровне кластера. Это приемлемый компромисс именно потому, что слот логической репликации Postgres остается вышестоящим надежным источником истины. Это wal2json, который гарантирует, что никакие изменения БД не будут потеряны перед использованием, а не NATS.

JetStream существует в aura-realtime, но для совершенно другого использования. Он обслуживает хранилище KV межэкземплярного присутствия (aura_presence, хранилище памяти, 60-секундное max_age), которое синхронизирует, кто и на каком канале находится в сети между всеми репликами ws-front. Кратковременное общее состояние, а не поток событий CDC, которые нужно воспроизвести.

#
Безопасность

РЛС перепроверяется абонентом на каждом мероприятии

Событие CDC не передается всем подписчикам канала. Опросчик направляет каждое изменение одному из рабочих процессов CDC_NUM_SHARDS (по умолчанию 4, хешируется project_id), который сначала опрашивает pg_policies. Таблица без политики RLS разрешает доступ всем подписчикам (поведение Supabase), тогда как таблица с политиками запускает проверку каждого подписчика.

subscriptions.rs::check_rls_visible_batch (реальный, упрощенный)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- JWT-заявления подписчика, внедренные в GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

Транзакция берет на себя физическую роль подписчика Postgres и внедряет свои утверждения JWT в request.jwt.claims — тот, который auth.uid() и auth.role() читают на стороне политики. Затем он проверяет, что строка остается видимой под этой ролью, а затем отменяет: запись не производится, настоящее чтение RLS. Только подтвержденный sub_id приземляется в allowed_sub_ids мероприятия; пустой vec означает никого, в отказоустойчивый.

Астуце

А DELETE обходит эту проверку построчно — строка исчезла, перечитать ее по RLS невозможно — и уходит обратно всем подписчикам канала. В свою очередь, полезные данные DELETE всегда предоставляют только столбцы идентификаторов (первичный ключ), но не содержимое удаленной строки.

Деталь, которая часто мешает миграции: wal2json захватывает полное содержимое UPDATE или DELETE только в том случае, если в таблице есть REPLICA IDENTITY FULL. Без него снова заработает только INSERT. Вот почему конечная точка, которая включает режим реального времени для таблицы, автоматически устанавливает этот ALTER TABLE — фактическое исправление репозитория, а не отдельный флажок.

#
Ловушка проверена

Почему CDC может хранить молчание в Kubernetes

История подачи документов свидетельствует о настоящей ловушке, а не о теоретическом случае. Образ Postgres, используемый по умолчанию в кластере Kubernetes (стандартный образ сообщества pgvector/pgvector), не включает плагин wal2json. Без него pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') терпит неудачу, и ничто на стороне клиента явно не сигнализирует об этом. WebSockets остаются открытыми, подписки принимаются, но никаких событий БД никогда не происходит.

docker/Postgres.Dockerfile (настоящий экстракт)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

Фактическое решение заключалось в том, чтобы создать и опубликовать специальный образ Postgres с установленным этим пакетом, а затем указать на него Kubernetes StatefulSet вместо стандартного образа. Он также явно установил wal_level=logical, max_replication_slots и max_slot_wal_keep_size в качестве аргументов запуска — они отсутствуют в общих образах.

Именно для того, чтобы сделать этот тип сбоя заметным, cdc-worker выставляет флаг cdc_active (ложный, пока слот не подтвержден как работоспособный) на своей выделенной конечной точке /health. Это двоичный сигнал, за которым следует следить, а не делать вывод о неработоспособности CDC только по отсутствию событий на стороне клиента.

#
Текущий предел

Что еще не охватывает этот конвейер: выделенные кластеры

Этот механизм считывает одну переменную POSTGRES_REPLICATION_URL, то есть один сервер Postgres и один логический слот. Это соответствует мультитенантной схеме для каждой модели проекта (project_<uuid>) в общем кластере. На этой модели все работает: единый cdc-worker видит изменения всех проектов в одном кластере и маршрутизирует по диаграмме.

Для проекта в выделенном кластере CNPG — собственном изолированном экземпляре Postgres — используемый сегодня образ (созданный на основе стандартного образа CloudNativePG) добавляет только pg_graphql, а не wal2json. Ничто также не запускает cdc-worker для каждого выделенного кластера. Таким образом, в текущем состоянии кода описанный здесь CDC реального времени остается сосредоточенным на общем уровне — реальной архитектурной границе, а не функциональности дорожной карты.

#
Клиентская часть

Что не меняется: API postgres_changes SDK

Весь этот механизм остается невидимым для SDK. Цепной APIaurabase-js не меняется в зависимости от того, поступает ли событие из cdc-worker через NATS или, в локальной разработке, из неразделенного двоичного файла aura-realtime.

app.tstypescript
const channel = aura.realtime.channel('room-1')
  .on('postgres_changes', {
    event: 'INSERT', schema: 'public', table: 'messages'
  }, (payload) => { /* … */ })
  .subscribe()

Но прежде чем произойдет одно событие, таблица должна быть зарегистрирована в CDC — это не происходит автоматически при создании. Вызов service_role на выделенной конечной точке (или эквивалентном переключателе в Studio) записывает таблицу в realtime_tables схемы платформы проекта и устанавливает для вас REPLICA IDENTITY FULL.

terminalbash
curl -X PUT \
  https://api.aurabase.cloud/v1/db/<project_id>/schema/tables/messages/realtime \
  -H "apikey: <service_role_key>" \
  -H 'Content-Type: application/json' \
  -d '{"enabled": true}'

Остальные сведения о Supabase-совместимом SDK — аутентификация, хранилище, политики RLS — см. в руководстве по миграции . О том, как aura-realtime разбивается на несколько двоичных файлов в одном контейнере рабочей области, см. в подробностях архитектуры Cargo.

ГОТОВЫ К РАЗВЕРТЫВАНИЮ?

Ваш бэкэнд за пять минут.

Кредитная карта не требуется · 500 МБ бесплатно · 50 000 MAU