Ответ
Apache Kafka — это распределенная потоковая платформа, идеально подходящая для обработки больших объемов данных в реальном времени. Она используется для построения масштабируемых и отказоустойчивых конвейеров данных.
Сценарий использования: Система аналитики пользовательских действий
Представьте веб-приложение, где необходимо отслеживать каждое действие пользователя (клики, просмотры страниц, добавление в корзину) для последующего анализа, персонализации или мониторинга.
Почему Kafka подходит для этой задачи:
- Масштабируемость: Легко обрабатывает миллионы событий в секунду.
- Надежность: Сохраняет данные на диске, обеспечивая отказоустойчивость и гарантированную доставку.
- Реальное время: Позволяет обрабатывать события практически мгновенно.
- Деcoupling: Отделяет генерацию событий от их обработки, позволяя добавлять новых потребителей без изменения продюсеров.
- Порядок сообщений: Гарантирует порядок сообщений в пределах одной партиции.
Пример кода (Python с kafka-python):
from kafka import KafkaProducer, KafkaConsumer
import json
from datetime import datetime
import time
# --- Producer (отправитель событий) ---
# Отправляет события о действиях пользователя в топик 'user_actions'
producer = KafkaProducer(
bootstrap_servers='localhost:9092', # Адрес Kafka брокера
value_serializer=lambda v: json.dumps(v).encode('utf-8') # Сериализация данных в JSON
)
def track_user_action(user_id: int, action: str):
event = {
'user_id': user_id,
'action': action,
'timestamp': datetime.now().isoformat()
}
producer.send('user_actions', event)
print(f"Отправлено событие: {event}")
# Пример отправки нескольких событий
# track_user_action(101, "page_view")
# track_user_action(101, "add_to_cart")
# track_user_action(102, "login")
# producer.flush() # Убедиться, что все сообщения отправлены
# --- Consumer (обработчик событий) ---
# Читает события из топика 'user_actions' и обрабатывает их
consumer = KafkaConsumer(
'user_actions',
bootstrap_servers='localhost:9092',
auto_offset_reset='earliest', # Начать чтение с самого начала топика при первом запуске
enable_auto_commit=True, # Автоматически коммитить оффсеты
group_id='analytics_group', # Группа потребителей
value_deserializer=lambda m: json.loads(m.decode('utf-8')) # Десериализация данных из JSON
)
print("Consumer запущен, ожидание сообщений...")
for message in consumer:
event_data = message.value
print(f"Получено событие: {event_data} из топика {message.topic}, партиции {message.partition}, оффсет {message.offset}")
# Здесь может быть логика обработки: сохранение в БД, отправка в другую систему, агрегация
# process_user_event(event_data)
time.sleep(0.1) # Имитация обработки
В этом примере:
- Producer отправляет структурированные события в топик
user_actions. - Consumer (или несколько consumer'ов) читает эти события, десериализует их и выполняет дальнейшую обработку (например, сохранение в аналитическую базу данных, обновление кэша или запуск других бизнес-процессов).
- Kafka выступает как буфер и надежный транспорт между источниками данных и их потребителями.