TGViewer
Data Apps Design Data Apps Design @data_apps · 2.08K subscribers
Post #394 1.56K
🔹 Хочу устроить Real Time Data Sync — дайте критику и комментарии

Привет! Последние пару недель плотно работаю над настройкой 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 | Навигация по каналу
  • ❤‍🔥 9
  • 🔥 6
  • ❤ 5
  • ⚡ 3
More from @data_apps
  1. Aug 25, 2026🔸 Меня заблокировал Cursor Сообщение: Your Cursor account was closed following an account…
  2. Feb 27, 2026✅ 3 ОФФЕРА, мои мысли и рекомендации по поиску работы в 2026 в Data и IT в целом Салют! Чу…
  3. Feb 6, 2026Эксперимент успешный 😌 Чек-лист: https://gist.github.com/kzzzr/e49b7e0b2af01e4e1dbc57102d…
  4. Feb 6, 2026👀 DataLens: Бесплатная сказка закончилась. Кейс миграции на Superset (Open Source BI) 1 м…
  5. Feb 2, 2026😘 Открываю доступ к закрытым записям цикла Designing Modern Data Apps Всем привет! 🟡 Это…
  6. Jan 22, 2026☄ Открыт к предложениям: Staff Data Engineer / Data Platform Lead За 11+ лет я прошел путь…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →