От теории переходим к практике — давай посмотрим, что реально происходит с данными.
Есть таблица
orders в PostgreSQL. Мы хотим держать её актуальную копию в аналитическом хранилище.🌺Наш пайплайн:
PostgreSQL → Debezium → Kafka → обработчик → хранилище.
PostgreSQL → Debezium → Kafka → обработчик → хранилище.
Сначала забираем то, что уже есть
В базе лежит заказ:
Debezium делает начальный снапшот, и обработчик переносит эту строку в хранилище. Теперь можно подхватывать изменения.
В базе лежит заказ:
id=101, status=new.Debezium делает начальный снапшот, и обработчик переносит эту строку в хранилище. Теперь можно подхватывать изменения.
1️⃣ Создали ещё один заказ
После фиксации изменения Debezium отправляет в Kafka событие создания строки. В нём есть ключ
Обработчик добавляет заказ в хранилище. Теперь там два заказа: 101 и 102.
INSERT INTO orders (id, status)
VALUES (102, 'new');
После фиксации изменения Debezium отправляет в Kafka событие создания строки. В нём есть ключ
102 и данные нового заказа.Обработчик добавляет заказ в хранилище. Теперь там два заказа: 101 и 102.
2️⃣ Первый заказ оплатили
Прилетает событие обновления. Обработчик находит строку по ключу
Заново выгружать всю таблицу не понадобилось 👍
UPDATE orders
SET status = 'paid'
WHERE id = 101;
Прилетает событие обновления. Обработчик находит строку по ключу
101 и меняет статус на paid.Заново выгружать всю таблицу не понадобилось 👍
🌿 Второй заказ удалили
Прилетает событие удаления с ключом
DELETE FROM orders
WHERE id = 102;
Прилетает событие удаления с ключом
102. Обработчик удаляет соответствующую строку из хранилища.После обработки всех событий в обеих системах остаётся один заказ: 101, paid.