Відзначимо: коли база даних 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 — і ми гарантуємо доставку змін за секунди.







