Собираем пайплайн: PostgreSQL → Debezium → Kafka → ClickHouse

Сквозная цепочка от изменения в PostgreSQL до строки в витрине: Debezium как коннектор Kafka Connect, outbox event router, snapshot и переход в streaming, сток в ClickHouse. Не про механику WAL и не про внутренности Kafka — про то, как это собирается вместе и где проходят швы

Опрашивать таблицу по updated_at в цикле — рабочий, но дорогой способ заметить, что в PostgreSQL что-то изменилось: лишняя нагрузка на базу, потерянные промежуточные состояния, невозможность поймать DELETE. Change data capture решает ту же задачу с другой стороны — читает журнал изменений СУБД и превращает каждый commit в поток событий. Как устроен этот журнал, чем логическая репликация отличается от физической и что такое pgoutput изнутри — тема отдельной статьи «WAL на службе: репликация, CDC и PITR»; здесь WAL — данность. Фокус этой статьи — на том, как из данности собирается пайплайн, который довозит изменение из PostgreSQL до строки в аналитической витрине ClickHouse: Debezium как коннектор Kafka Connect, таблица outbox как единственный источник, Kafka как транспорт, явный консьюмер как приёмник.

Это первая статья серии «Стриминг и обработка данных». Транспорт (Kafka и выбор брокера вообще) уже разобран в обзорной статье про выбор брокера — здесь Kafka уже данность, а Kafka Connect как раннер коннекторов подробнее показан в «Экосистема Kafka: Connect, Schema Registry, ksqlDB». За текстом статьи — стенд digital-cookbook/messaging/data-streaming: PostgreSQL с таблицей outbox, Kafka Connect с Debezium PostgreSQL connector, топик orders.events, Kafka Streams, сворачивающий поток в агрегат customer.totals, и Go-консьюмер, укладывающий агрегаты в ClickHouse. Все числа ниже — из живого прогона этого стенда.

Ретрофутуристская сборочная линия в духе «Полдня»: слева бак-хранилище PostgreSQL, из него по прозрачной трубе течёт лента WAL; механический глаз-коннектор Debezium считывает с ленты только карточки с меткой outbox, штампует их и подаёт на конвейер-жёлоб Kafka; в конце жёлоба ковш-консьюмер зачерпывает карточки рядами и укладывает в картотеку ClickHouse. Рядом — второй, пунктирный жёлоб напрямую из Kafka в картотеку, перечёркнутый, с табличкой «дедуп — всё равно на вас»

В статье

Карта цепочки: что где живёт

Прежде чем разбирать конфиги, полезно увидеть весь путь одного изменения целиком — от COMMIT в PostgreSQL до строки в ClickHouse:

Звено Роль
PostgreSQL, таблица outbox единственная таблица, которую видит коннектор; домен (orders, customers и т. д.) в неё не входит
Kafka Connect + Debezium PostgreSQL connector читает WAL через pgoutput, превращает commit-ы в события
Outbox Event Router (SMT) разворачивает строку outbox в доменное событие и решает, в какой топик её положить
Kafka, топик orders.events транспорт сырых доменных событий; ключ сообщения — aggregateid, что определяет партицию и порядок
Kafka Streams сворачивает поток orders.events в топик customer.totals — поток последовательных снимков: новое накопительное состояние клиента после каждого заказа, актуальна последняя версия; state store, changelog и гарантии обработки — вторая статья серии
Kafka, топик customer.totals транспорт готового агрегата — то, что фактически читает сток
Go-консьюмер читает customer.totals, пишет строки в ClickHouse явным кодом, а не через встроенный движок
ClickHouse аналитическая витрина, конечная точка цепочки

Каждое звено решает одну узкую задачу и не лезет в соседнюю: коннектор не знает о ClickHouse, Kafka Streams не знает про ClickHouse, консьюмер не знает, как устроен WAL. Дальше — по порядку, начиная с того, что видит сам коннектор.

Debezium в Kafka Connect: конфиг и что подключено

Debezium — не отдельный сервис, а плагин для Kafka Connect: JVM-раннер, который поднимает и конфигурирует коннекторы через REST API, без единой строки кода на потребителя. Сам Connect и его экосистема (Schema Registry, SMT-цепочки, distributed-режим) разобраны отдельно в статье про экосистему Kafka; здесь — только то, что нужно для этого конкретного коннектора.

Регистрация коннектора — один JSON, отправленный на REST worker’а:

{
  "name": "ds-outbox-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "shop",
    "database.password": "shop",
    "database.dbname": "shop",
    "topic.prefix": "ds",
    "plugin.name": "pgoutput",
    "publication.name": "dbz_pub",
    "slot.name": "ds_slot",
    "table.include.list": "public.outbox",
    "snapshot.mode": "initial",
    "transforms": "outbox",
    "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
    "transforms.outbox.route.by.field": "aggregatetype",
    "transforms.outbox.route.topic.replacement": "${routedByValue}.events",
    "transforms.outbox.table.field.event.key": "aggregateid",
    "transforms.outbox.table.field.event.payload": "payload",
    "transforms.outbox.table.expand.json.payload": "true",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter"
  }
}

Три поля стоит выделить сразу:

  • plugin.name: pgoutput — встроенный (с PostgreSQL 10) output-плагин логической репликации. Как именно он декодирует WAL в поток изменений — за пределами этой статьи, механика в статье про WAL. Здесь важно только то, что pgoutput не требует установки дополнительных расширений — коннектор работает с чистым образом postgres:18.4.
  • table.include.list: public.outbox — коннектор физически не видит ничего, кроме outbox. Это не оптимизация, а осознанная граница: домен (orders, customers, их схема, внутренние ключи) наружу не течёт, коннектор подписан только на контракт. Почему этот контракт — отдельная таблица outbox, а не CDC по доменным таблицам напрямую, разбирает статья про transactional outbox/inbox — паттерн, на котором держится вся эта схема.
  • snapshot.mode: initial — при первой регистрации коннектор сначала выгружает текущее содержимое outbox целиком (snapshot), потом переключается на поток изменений (streaming). Что из этого получилось на живом стенде — ниже.

Snapshot и переход в streaming

Коннектор был зарегистрирован на только что созданную таблицу outbox — пустую. Дословный лог docker logs ds-connect сразу после регистрации:

2026-07-17T18:25:19,340 INFO   Postgres||snapshot  Snapshot step 7 - Snapshotting data   [io.debezium.relational.RelationalSnapshotChangeEventSource]
2026-07-17T18:25:19,351 INFO   Postgres||snapshot  	 Finished exporting 0 records for table 'public.outbox' (1 of 1 tables); total duration '00:00:00.005'   [io.debezium.relational.RelationalSnapshotChangeEventSource]
2026-07-17T18:25:19,352 INFO   Postgres||snapshot  Snapshot completed   [io.debezium.pipeline.source.AbstractSnapshotChangeEventSource]
2026-07-17T18:25:19,353 INFO   Postgres||snapshot  Snapshot ended with SnapshotResult [status=COMPLETED, ...]   [io.debezium.pipeline.ChangeEventSourceCoordinator]
2026-07-17T18:25:19,375 INFO   Postgres||streaming  Starting streaming   [io.debezium.pipeline.ChangeEventSourceCoordinator]
2026-07-17T18:25:19,390 INFO   Postgres||streaming  Starting replication stream from LSN LSN{0/1BB7648} with automaticFlush=false (mode=CONNECTOR)   [io.debezium.connector.postgresql.connection.PostgresReplicationConnection]
2026-07-17T18:25:19,404 INFO   Postgres||streaming  Processing messages   [io.debezium.connector.postgresql.PostgresStreamingChangeEventSource]

Finished exporting 0 records — снапшот снял 0 строк, потому что таблица outbox была пуста на момент регистрации. Уже через несколько миллисекунд коннектор переключился в режим streaming и встал в Processing messages.

Дальше запущен bash scripts/seed.sh 50 — 50 заказов, каждый одной транзакцией пишущий строку в outbox. Снапшот их захватить не мог: он уже завершился до того, как эти строки появились в PostgreSQL. Значит все 50 событий обязаны были пройти именно через streaming. Debezium на уровне INFO не логирует каждое стримингованное событие отдельно (нормальное поведение pgoutput, подробности — за пределами этой статьи), поэтому доказательство — не строка лога, а фактическое число сообщений в топике:

$ docker exec ds-kafka /opt/kafka/bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic orders.events
orders.events:0:22
orders.events:1:17
orders.events:2:11

22 + 17 + 11 = 50 — ровно столько строк вставил seed.sh. Ни одна не могла попасть через snapshot (он закончился раньше, чем строки появились), значит все 50 прошли через streaming. Это и есть переход snapshot → streaming, наблюдаемый не по описанию, а по цифрам, которые обязаны сойтись.

Outbox event router: почему тело события «голое»

Строка в outbox — это техническая обёртка: id, aggregatetype, aggregateid, type, payload, created_at. Наружу в Kafka должно уйти не это, а доменное событие. Разворачивает обёртку в событие Outbox Event Router — SMT (single message transform), сконфигурированная прямо в JSON коннектора четырьмя полями:

  • route.by.field: aggregatetype + route.topic.replacement: ${routedByValue}.events — топик назначения вычисляется из значения колонки aggregatetype конкретной строки. Отсюда и имя топика в примерах выше — orders.events: значит, aggregatetype для заказов в этом стенде равен orders.
  • table.field.event.key: aggregateid — ключом сообщения в Kafka становится aggregateid, а не id строки outbox. Это определяет партицию и, следовательно, порядок доставки — почему это важно, разбирает третья статья серии (ниже, в разделе про границы).
  • table.field.event.payload: payload + expand.json.payload: true — значение jsonb-колонки payload не оборачивается во внешний конверт, а разворачивается и становится телом сообщения напрямую. Вместе с value.converter.schemas.enable: false (без envelope-схемы JsonConverter) итоговое сообщение в Kafka — это ровно тот JSON, что лежал в payload, и ничего больше.

Отсюда и «голое» тело: событие в orders.events не нужно доставать из вложенной структуры вида {"payload": {...}, "schema": {...}} — оно уже плоское. Консьюмеру на другом конце не нужен двойной parse — распаковка сделана роутером один раз, на входе в Kafka, а не на каждом чтении.

Почему вообще таблица-посредник, а не CDC прямо по доменным таблицам (orders, customers) — отдельный вопрос дизайна, который решает сам паттерн: Transactional Outbox/Inbox держит формат события под контролем автора домена и не даёт внутренней схеме БД утечь в контракт для внешних потребителей. Здесь этот паттерн только используется, не разбирается заново.

Сток: явный консьюмер вместо Kafka engine

Событие долетело до Kafka плоским JSON-ом — но прежде чем оно окажется строкой в ClickHouse, оно проходит ещё через одно звено. orders.events — это сырые события, по одному на заказ; в ClickHouse нужна не история заказов, а актуальная сумма по клиенту. Между ними стоит Kafka Streams: приложение сворачивает поток orders.events в топик customer.totals, куда после каждого заказа кладётся новый снимок текущей суммы клиента — то есть на клиента приходится много версий, из которых актуальна последняя (как из этого потока снимков читается именно последнее состояние на стоке — разбирает третья статья серии). Go-консьюмер в стенде читает именно customer.totals, а не сырой orders.events напрямую: парсит JSON и вставляет строки в ClickHouse-таблицы явным кодом. Как устроена сама агрегация — state store, changelog, восстановление после падения, гарантии exactly-once — за пределами этой статьи, это тема второй статьи серии, «Kafka Streams: обработка без кластера». Здесь важно только то, что консьюмеру и ClickHouse достаётся уже готовый агрегат, а не сырой поток заказов, который пришлось бы сворачивать самостоятельно на каждой вставке.

Почему не проще. ClickHouse умеет читать Kafka сам — движок ENGINE=Kafka подписывается на топик и материализует данные через материализованное представление; устройство этого пути и когда его стоит брать разбирает статья про материализованные представления и real-time агрегации в ClickHouse. Но сам движок ENGINE=Kafka дедупликацию не даёт — он избавляет только от кода консьюмера: poll/parse/insert за вас делает ClickHouse. Вопрос дублей всё равно приходится решать — ReplacingMergeTree на целевой таблице, дедуп в материализованном представлении или гарантии upstream, — просто в другом месте схемы, а не в коде явного консьюмера. Здесь мы берём явный консьюмер, чтобы это место — точку, где принимается решение о дублях и идемпотентности, — было видно на карте цепочки отдельным звеном, а не растворялось в конфигурации движка хранилища.

Kafka по умолчанию доставляет at-least-once: при ретраях консьюмера дубли на входе в ClickHouse неизбежны в принципе, вне зависимости от того, кто читает топик — свой код или встроенный движок. Что конкретно делать с этими дублями (идемпотентный consumer, ключ дедупликации, ReplacingMergeTree) — тема отдельной статьи серии, ссылка ниже. Здесь важно только то, что решение принято осознанно: явный консьюмер как отдельное, видимое звено цепочки, а не скрытая часть движка хранилища.

Границы: что не разобрано здесь

Эта статья — про сборку, а не про механику внутри каждого звена. За разбором того, что осталось за кадром, — по ссылкам:

Демо и версии

Стенд — digital-cookbook/messaging/data-streaming: PostgreSQL с таблицей outbox, Kafka + Kafka Connect с зарегистрированным Debezium PostgreSQL connector и Outbox Event Router, Kafka Streams, сворачивающий orders.events в customer.totals, Go-консьюмер, пишущий в ClickHouse. docker-compose, скрипт seed.sh и шаги воспроизведения снапшота/streaming — в README демо.

Версии на живом прогоне: PostgreSQL 18.4, Kafka 4.3.1, Kafka Connect (framework) 4.3.0, Debezium connector 3.6.0.Final, ClickHouse 26.6.1.1193. Поведение snapshot.mode, Outbox Event Router и совместимость коннектора с версией Kafka зафиксированы на этом наборе — при обновлении любого компонента стоит перепроверить.

Документация

Обсуждение в Telegram

Присоединиться →

Комментарии