Налаштування Kafka Connect для інтеграції з базами даних

Відзначимо: коли база даних PostgreSQL розростається до сотень гігабайт, а вимоги до актуальності пошукового індексу — секунди, ручна синхронізація перестає працювати. Ми налаштовуємо Kafka Connect для потокової реплікації даних (CDC) — це надійніше та швидше, ніж кастомні Poller-сервіси. Наприклад,

Розробка та обслуговування будь-яких видів сайтів:

Інформаційні сайти або веб-програми
Сайти візитки, landing page, корпоративні сайти, онлайн каталоги, квіз, промо-сайти, блоги, ресурси новин, інформаційні портали, форуми, агрегатори
Сайти або веб-програми електронної комерції
Інтернет-магазини, B2B-портали, маркетплейси, онлайн-обмінники, кешбек-сайти, біржі, дропшиппінг-платформи, парсери товарів
Веб-програми для управління бізнес-процесами
CRM-системи, ERP-системи, корпоративні портали, системи управління виробництвом, парсери інформації
Сайти або веб-програми електронних послуг
Дошки оголошень, онлайн-школи, онлайн-кінотеатри, конструктори сайтів, портали надання електронних послуг, відеохостинги, тематичні портали

Це лише деякі з технічних типів сайтів, з якими ми працюємо, і кожен із них може мати свої специфічні особливості та функціональність, а також бути адаптованим під конкретні потреби та цілі клієнта.

Послуги, які ми пропонуємо
Показано 1 з 1Усі 2062 послуг
Налаштування Kafka Connect для інтеграції з базами даних
Складний
~3-5 днів

Наші компетенції:

Часті запитання

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1419
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1287
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    983
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1245
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    983
  • image_bitrix-bitrix-24-1c_fixper_448_0.webp
    Розробка веб-сайту для компанії ФІКСПЕР
    998

Відзначимо: коли база даних PostgreSQL розростається до сотень гігабайт, а вимоги до актуальності пошукового індексу — секунди, ручна синхронізація перестає працювати. Ми налаштовуємо Kafka Connect для потокової реплікації даних (CDC) — це надійніше та швидше, ніж кастомні Poller-сервіси. Наприклад, маркетплейс з каталогом у PostgreSQL та пошуком в Elasticsearch стикається з затримками оновлення індексу до 15 хвилин. З нашим рішенням затримка скорочується до 2 секунд, а при збої вузла дані не втрачаються — таски перерозподіляються автоматично. Один з типових кейсів: PostgreSQL → Debezium → Kafka → Elasticsearch. Без Kafka Connect інженери витрачають тижні на написання WAL-обробника та боротьбу з дублікатами. Ми робимо це за 4 дні під ключ, з відмовостійким кластером та моніторингом. Документація Debezium підтверджує, що CDC забезпечує потокову передачу з мінімальною затримкою.

Чому Kafka Connect краще кастомної інтеграції?

Порівняємо підходи:

Критерій Кастомний сервіс Kafka Connect + Debezium
Час розробки 2–4 тижні 4 дні
Дублікати при збоях Так, потрібна ідемпотентність Автоматично, завдяки offset-ам
Моніторинг Свій код метрик Вбудований REST API + Prometheus
Масштабування Ручне, з переписуванням Distributed-режим, додавання вузлів
Підтримка типів даних Для кожного типу — своя серіалізація Avro/JsonSchema/Protobuf через Schema Registry
Відмовостійкість Ручна, рестарт сервісу Автоматичний ребаланс тасків

Результат: Kafka Connect дає готовий фреймворк з гарантованою доставкою (exactly-once з ідемпотентними продюсерами) та економить 80% часу на розробку і 50% на експлуатацію.

Як працює CDC на PostgreSQL з Debezium?

Change Data Capture (CDC) перехоплює кожну зміну в базі даних та транслює її в подієвий потік. Debezium підключається до WAL (Write-Ahead Log) PostgreSQL та відправляє INSERT/UPDATE/DELETE в Kafka. Це позбавляє необхідності писати власні тригери або опитувати таблиці за розкладом. Debezium підтримує режим snapshot.mode=initial для початкового завантаження та incremental для уникнення блокувань на великих таблицях. Після налаштування кожна зміна з'являється в Kafka-топіку за мілісекунди. Середній лаг становить 100 мс, пропускна здатність — до 10000 повідомлень/сек.

Налаштування Kafka Connect під ключ: етапи

Процес складається з 4 етапів, кожен з перевіркою якості.

Етап 1: Підготовка PostgreSQL та Kafka

Вмикаємо логічну реплікацію: wal_level = logical, створюємо publication для потрібних таблиць та користувача debezium з правами SELECT. На стороні Kafka перевіряємо bootstrap.servers, налаштування ретеншену та compact-топіки для Debezium.

Етап 2: Розгортання Kafka Connect в distributed-режимі

Кластер з 2–3 вузлів з внутрішніми топіками для конфігурації. Конфігурація повністю типізована (див. нижче). Використовуємо Avro з Schema Registry — це дає гарантію сумісності схем при зміні бази.

Етап 3: Налаштування Debezium Source Connector

Debezium читає WAL PostgreSQL та відправляє кожну зміну в Kafka. Налаштовуємо snapshot.mode=initial, transforms для витягу нового рядка та tombstone для DELETE. Для великих таблиць (мільярди рядків) використовуємо incremental snapshot, щоб не блокувати БД.

Етап 4: Sink-конектори в Elasticsearch та PostgreSQL

Для пошукового індексу — Elasticsearch Sink з батчінгом 500 записів та retry backoff. Для аналітичної БД — JDBC Sink з upsert та pk.mode=record_key. На кожному етапі тестуємо INSERT/UPDATE/DELETE та лаг.

Як розгорнути Kafka Connect та конектори?

Нижче — ключові конфігурації для запуску distributed-кластера та типових конекторів.

# distributed-властивості (connect-distributed.properties) bootstrap.servers=kafka-1:9092,kafka-2:9092,kafka-3:9092 group.id=kafka-connect-cluster config.storage.topic=connect-configs offset.storage.topic=connect-offsets status.storage.topic=connect-statuses config.storage.replication.factor=3 offset.storage.replication.factor=3 status.storage.replication.factor=3 offset.flush.interval.ms=10000 rest.host.name=0.0.0.0 rest.port=8083 rest.advertised.host.name=connect-1.internal rest.advertised.port=8083 plugin.path=/opt/kafka/plugins key.converter=io.confluent.connect.avro.AvroConverter key.converter.schema.registry.url=http://schema-registry:8081 value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://schema-registry:8081 

Перед запуском Debezium налаштуйте PostgreSQL:

ALTER SYSTEM SET wal_level = logical; ALTER SYSTEM SET max_replication_slots = 10; ALTER SYSTEM SET max_wal_senders = 10; CREATE USER debezium WITH REPLICATION LOGIN PASSWORD 'secure_password'; GRANT CONNECT ON DATABASE myapp TO debezium; GRANT USAGE ON SCHEMA public TO debezium; GRANT SELECT ON ALL TABLES IN SCHEMA public TO debezium; ALTER DEFAULT PRIVILEGES IN SCHEMA public GRANT SELECT ON TABLES TO debezium; CREATE PUBLICATION debezium_pub FOR TABLE products, orders, users, categories; 

Тепер зареєструйте конектори через REST API. Ось приклад Debezium Source та JDBC Sink в одному блоці:

# Debezium Source curl -X POST http://connect-1:8083/connectors -H "Content-Type: application/json" -d '{ "name": "postgres-source-connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.hostname": "postgres.internal", "database.port": "5432", "database.user": "debezium", "database.password": "secure_password", "database.dbname": "myapp", "database.server.name": "myapp-pg", "topic.prefix": "myapp", "table.include.list": "public.products,public.orders,public.users", "plugin.name": "pgoutput", "publication.name": "debezium_pub", "slot.name": "debezium_slot", "snapshot.mode": "initial", "snapshot.isolation.mode": "read_committed", "decimal.handling.mode": "double", "time.precision.mode": "connect", "tombstones.on.delete": "true", "heartbeat.interval.ms": "10000", "transforms": "unwrap", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.delete.handling.mode": "rewrite", "transforms.unwrap.add.fields": "op,ts_ms,source.ts_ms" } }' # JDBC Sink curl -X POST http://connect-1:8083/connectors -H "Content-Type: application/json" -d '{ "name": "postgres-sink-connector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "4", "topics": "myapp.analytics.events", "connection.url": "jdbc:postgresql://analytics-pg:5432/analytics", "connection.user": "kafka_writer", "connection.password": "secure_password", "auto.create": "false", "auto.evolve": "false", "insert.mode": "upsert", "pk.mode": "record_key", "pk.fields": "id", "table.name.format": "analytics.${topic}", "batch.size": "1000", "db.timezone": "UTC", "transforms": "dropPrefix", "transforms.dropPrefix.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.dropPrefix.exclude": "__deleted,__op,__ts_ms" } }' 

Управління конекторами виконується через REST API: отримання статусу, пауза, перезапуск тасків, що впали. Prometheus JMX-метрики налаштовуються через JMX Exporter.

Типові проблеми та їх вирішення

WAL bloatЯкщо слот реплікації не зсувається, WAL накопичується. Налаштовуємо `max_slot_wal_keep_size` у PostgreSQL та алерт на розмір WAL. Регулярно моніторимо та чистимо.
Schema evolutionПри додаванні нової колонки Debezium автоматично оновить схему в Schema Registry. Sink-конектор повинен бути готовий (auto.evolve=true або ручне керування).
Tombstone messagesПри DELETE Debezium відправляє два повідомлення: подію DELETE та tombstone (null value). Для compact-топіків tombstone видаляє запис з логу.

Що входить у роботу

  • Аналіз поточної схеми БД та навантажень, вибір конекторів та трансформацій
  • Налаштування PostgreSQL для логічної реплікації (WAL, publication, користувачі)
  • Встановлення Kafka Connect у distributed-режимі на 2–3 вузли з Schema Registry
  • Розгортання Debezium Source Connector з initial snapshot
  • Налаштування одного або декількох Sink-конекторів (Elasticsearch, JDBC, S3)
  • Написання Single Message Transforms (SMT) для підгонки схеми
  • Інтеграція моніторингу: Prometheus JMX Exporter, дашборд Grafana, алерти в Slack
  • Документація схеми топіків та конфігурацій
  • Навчання команди (2 години: базові операції, рестарт, діагностика)

Таймлайн

День Робота
1 Налаштування PostgreSQL для логічної реплікації, встановлення Kafka Connect у distributed-режимі на 2–3 вузли
2 Встановлення Debezium, первинний snapshot (може зайняти години для великих таблиць), налаштування конектора, верифікація CDC-подій
3 Налаштування Sink-конектора (ES або PostgreSQL), трансформації через SMT, тестування повного пайплайну INSERT/UPDATE/DELETE
4 Моніторинг, алерти на лаг та помилки, документація схеми топіків, навантажувальне тестування з піковим потоком змін

За 10 років ми реалізували 50+ інтеграційних пайплайнів на PostgreSQL, MySQL та MongoDB. Отримайте консультацію щодо вашого проєкту – наші інженери проаналізують схему та навантаження і запропонують оптимальну архітектуру. Зв'яжіться з нами — оцінимо ваш проєкт за 2 години безкоштовно. Замовте налаштування Kafka Connect — і ми гарантуємо доставку змін за секунди.