Привет! Последние пару недель плотно работаю над настройкой Data Pipelines через Debezium + Kafka Connect.
🔺 Общая схема будет выглядеть так:
— Sources: MongoDB, Postgres
— Source connectors: Kafka Connect + Debezium (CDC)
— Kafka as intermediate storage
— Sink connector: Snowflake connector for Kafka (Snowpipe streaming)
— Snowflake as destination
Упражнения, которые я проделываю:
🟡 Deploy services: Kafka + Kafka Connect (+ Zookeeper + Schema registry + REST Proxy)
Вообще, готовые образы есть от: Confluent / Debezium / Strimzi.
Я пока использую образы от Debezium + Strimzi. Модифицирую, добавляя нужные мне коннекторы указанных версий (JAR-файлы).
У Strimzi есть операторы для K8s, которые я планирую использовать.
А вот K8s deployment от Confluent входит в Enterprise plan и подлежит лицензированию.
🔵 Choose collections for replication
Выбрал несколько несложных коллекций для тестирования и дебага.
Несложные, потому что:
— Небольшое количество колонок
— Минимальное количество вложеных nested-структур
— При этом частые обновления - получаю лог изменений CDC сразу.
— Глазами можно отследить все применяемые изменения (либо то, что не получается)
🟢 Source connector (MongoDB)
На выбор у меня было 2 Connectors:
— MongoDB Kafka Connector
— Debezium connector for MongoDB
Несмотря на наличие отличной документации и даже поддержку snapshot с учетом фильтра (например, sync истории только за 1 год), выбор сделан в пользу Debezium Connector. Позже подробно напишу доступные возможности и пример конфигурации.
🩷 Metadata fields for events
В каждое событие дополнительно пишу ряд метаданных:
—
ts_ms = metadata timestamp—
op = c (create), u (update), d (delete), r (snapshot)🔴 Message transformations: Topic route, New Document State Extraction
— Topic route: Мессаджи направляю сразу в топики с целевыми названиями таблиц в Snowflake.
— New Document State Extraction: этат трансформация позволяет перейти от сложного и детального формата события Debezium (schema, before, after) к упрощенному виду (только after), который мне и нужен в Snowflake.
🟤 Debug events
Развернул devcontainer (docker-compose based), установил kcat + jq.
Использую kcat для просмотра топиков, key-value значений, оффсетов.
jq использую для форматирования ответов от REST API (Kafka Connect).
Очень удобно проконтролировать изменения, которые вношу пунктом ранее.
🟡В следующем посте:
— Externalizing secrets
— Create a sink connector (Snowflake)
— Set up monitoring (Prometheus + Grafana)
— Use AVRO data format (
io.confluent.connect.avro.AvroConverter)— Configure topics (compact + delete, etc.)
— Configure incrmental snapshots loads via Debezium signals
— Ensure PII data masking (message transformations, connector configuration)
— Add HTTP bridge / REST proxy
🔻 Могли бы покритиковать / поделиться опытом / посоветовать что-либо?
🌐 @data_apps | Навигация по каналу