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







