Агрегация данных во времени с Kafka Streams
О чём речь
В FunBox мы делаем продукты для мобильных операторов: порталы, геосервисы, платежи, мобильную рекламу и многое другое. Один из проектов построен на микросервисной архитектуре, а основная функциональность связана с обработкой потоков событий.
Для организации централизованного, масштабируемого и быстрого обмена сообщениями используется Apache Kafka. В сервисе аналитики звонков возникла задача агрегировать несколько потоков данных во времени.
{"user1": id1, "user2": id2, "timestamp": ts1, "analytics": data1}
{"user1": id1, "user2": id2, "timestamp": ts2, "metrics": data2}Приложение получает события в соответствующих топиках Kafka. Время событий может не совпадать, но различие допускается не более чем на заданную константу T.
Основные понятия Kafka Streams
Kafka Streams — клиентская библиотека для потоковой работы с данными в Kafka. Она предоставляет высокоуровневый Streams DSL и низкоуровневый Processor API. Центральное понятие — топология сети обработчиков: граф потоковых процессоров и потоков.
Решение включает получение и группировку входящих событий по паре абонентов, объединение событий одного звонка и отсеивание промежуточных записей.
Комментарии (0):
К публикации пока нет комментариев.