⚙️ Apache Kafka: экосистема
Apache Kafka - платформа для потоковых данных
Включает:
⚫️Kafka Connect
⚫️ksqlDB
⚫️Kafka Streams
Kafka Streams
Библиотека для обработки потоков событий с возможностью:
🤩 Агрегации: подсчет количества событий за период для каждого ключа
➡️ количества кликов по рекламе для каждого пользователя за последний час
🤩Обогащении: дополнение событий данными из внешних систем или других топиков
➡️ добавление информации о профиле пользователя к событиям покупок
🤩Фильтрации: отбор нужных событий
🤩Трансформации: изменение формата/структуры сообщения
➡️ конвертация из бинарного формата Avro в JSON
🤩Объединения: данные из нескольких топиков
➡️ объединение данных о заказах и платежах
Архитектурная идея: микросервисный подход к потоковой обработке
Kafka Streams инкапсулирует логику обработки в независимое приложение, оно масштабируется вместе с кластером Kafka
🤩Обадает отказоустойчивостью
🤩Не требует развертывания отдельной инфраструктуры
Хранение данных между обработками
Для этого используется:
✨State Store — локальное хранилище внутри Kafka Streams, где находятся текущие вычисления
✨Changelog Topic — специальный топик в Kafka, куда записываются изменения в State Store
Если приложение перезапускается, то загружает данные из этого топика и продолжает работу с того же места
По умолчанию Kafka Streams хранит состояние локально, но обработанные данные можно записывать во внешние БД или облачные хранилища
Примеры использования
🤩Обработка транзакций с добавлением информации о пользователе из внешней БД
🤩 обогащение данных о платеже информацией о возрастной группе и истории покупок пользователя для системы фрод-мониторинга.
🤩Трансформация данных из формата Avro в JSON с валидацией и фильтрацией некорректных записей
🤩 очистка и преобразование данных логов веб-сервера перед загрузкой в аналитическое хранилище
Kafka Connect
Фреймворк для масштабируемого ввода/вывода данных между Kafka и внешними системами
Решает задачи интеграции с различными источниками
⏩ Пример: синхронизация данных между PostgreSQL и Elasticsearch
Источник (JDBC Connector) читает изменения из БД с помощью механизма изменения данных Debezium, а приемник (Elasticsearch Connector) загружает данные в поисковый индекс для быстрого поиска
Особенности
⏺готовые коннекторы для популярных систем
⏺автоматическое управление смещениями: отслеживание позиции обработки для каждого коннектора
⏺масштабирование через распределенный режим работы в кластере
⏺поддержка преобразований данных: встроенные онлайн преобразования форматов данных
ksqlDB
СУБД для потоковой обработки
➖позволяет выполнять SQL-запросы к данным в топиках Kafka
➖применяется для быстрого прототипирования и простых ETL-задач без кода на Java
🤩Пример: мониторинг аномальной активности пользователей в реальном времени.
Например, обнаружение пользователей, выполняющих более 100 действий в минуту, с отправкой уведомлений в систему безопасности
Особенности
🤩SQL-синтаксис для потоковой обработки
🤩поддержка оконных агрегаций и joins - возможность объединения потоков данных и агрегации по временным окнам
🤩REST API для управления потоковыми запросами
Применение: создание реальных дашбордов, реализация простых правил бизнес-логики, мониторинг качества данных
📎 Материалы
1. Kafka Streams (official site)
2. Экосистема Apache Kafka: Kafka Streams, Kafka Connect
3. Потоковая обработка данных с помощью Kafka Streams: архитектура и ключевые концепции
4. Под капотом Kafka Connect: источники, приемники и коннекторы
5. ksqlDB
6. ksqlDb или SQL как инструмент обработки потоков данных
📚 Книги
1. Kafka Streams и ksqlDB: данные в реальном времени - Сеймур Митч
2. Kafka в действии - Дилан Скотт, Виктор Гамов и Дейв Клейн (Глава 12)
#интеграции
➿➿➿➿➿➿➿➿
🧑🎓 Больше полезного в базе знаний по системному анализу
Post #637
15K
- ❤ 13
- 👍 7
- 🔥 5
- ⚡ 1