Опрашивать таблицу по 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. Все числа ниже — из живого прогона этого стенда.
В статье
- Карта цепочки: что где живёт
- Debezium в Kafka Connect: конфиг и что подключено
- Snapshot и переход в streaming
- Outbox event router: почему тело события «голое»
- Сток: явный консьюмер вместо Kafka engine
- Границы: что не разобрано здесь
- Демо и версии
- Документация
Карта цепочки: что где живёт
Прежде чем разбирать конфиги, полезно увидеть весь путь одного изменения целиком — от 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) — тема отдельной статьи серии, ссылка ниже. Здесь важно только то, что решение принято осознанно: явный консьюмер как отдельное, видимое звено цепочки, а не скрытая часть движка хранилища.
Границы: что не разобрано здесь
Эта статья — про сборку, а не про механику внутри каждого звена. За разбором того, что осталось за кадром, — по ссылкам:
- Механика WAL и логической репликации (слоты,
REPLICA IDENTITY, output-плагины) — «WAL на службе: репликация, CDC и PITR». - Сам паттерн outbox (relay/poller, inbox на стороне потребителя, idempotency key) — «Transactional Outbox/Inbox».
- Kafka Connect и его экосистема целиком (distributed worker, Schema Registry, ksqlDB) — «Экосистема Kafka: Connect, Schema Registry, ksqlDB».
- Stateful-обработка потока между сырыми событиями и готовыми агрегатами (Kafka Streams как библиотека внутри сервиса) — вторая статья серии, «Kafka Streams: обработка без кластера».
- Идемпотентные консьюмеры, порядок, DLQ и повторы — третья статья серии, «Сквозные гарантии и жизнь пайплайна».
- Kafka Streams и Flink изнутри (state stores, exactly-once, event-time, watermarks) — отдельная серия «Apache Flink: глубокое погружение».
Демо и версии
Стенд — 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 зафиксированы на этом наборе — при обновлении любого компонента стоит перепроверить.
Документация
- Debezium — документация и Debezium PostgreSQL connector — параметры коннектора,
snapshot.mode, режимы output-плагинов. - Debezium Outbox Event Router — полный список полей SMT и режимов роутинга.
- Kafka Connect — документация — REST API, distributed worker, конфигурация коннекторов.
- Смежное на сайте: хаб «Messaging: выбрать и эксплуатировать», обзор выбора брокера, «WAL на службе: репликация, CDC и PITR», «Transactional Outbox/Inbox».
Комментарии