Ответ
Apache Kafka — это распределенная стриминговая платформа, работающая по модели «издатель-подписчик» (pub/sub). В ее основе лежит концепция распределенного, неизменяемого, реплицируемого лога коммитов.
Основные компоненты:
- Broker (Брокер): Сервер Kafka, который хранит данные. Кластер Kafka состоит из нескольких брокеров для обеспечения отказоустойчивости и масштабируемости.
- Topic (Топик): Именованный поток (категория) сообщений. Например,
user-registrationsилиorder-updates. - Partition (Партиция): Топик делится на одну или несколько партиций. Партиция — это упорядоченный, неизменяемый лог сообщений. Это основная единица параллелизма в Kafka.
- Offset (Смещение): Уникальный последовательный номер, который идентифицирует каждое сообщение внутри партиции.
- Producer (Издатель): Приложение, которое отправляет (публикует) сообщения в топики Kafka.
- Consumer (Подписчик): Приложение, которое читает (потребляет) сообщения из топиков.
- Consumer Group (Группа подписчиков): Группа подписчиков, которые совместно читают сообщения из одного или нескольких топиков. Kafka гарантирует, что каждая партиция будет обрабатываться только одним подписчиком из группы, что позволяет легко распараллеливать обработку.
- ZooKeeper / KRaft: Внешний сервис (ZooKeeper) или встроенный протокол (KRaft в новых версиях) для управления метаданными кластера: информацией о брокерах, топиках, партициях и правах доступа.
Ключевые принципы работы:
- Неизменяемый лог: Сообщения записываются в конец партиции и не могут быть изменены. Они удаляются только по истечении срока хранения (retention policy).
- Распределенность и репликация: Партиции распределяются по брокерам в кластере. Каждая партиция имеет одну
leader-реплику (для чтения и записи) и несколькоfollower-реплик (для отказоустойчивости). - Состояние у клиента: В отличие от старых брокеров, Kafka не отслеживает, какие сообщения были прочитаны. Эту информацию (offset) хранит сам Consumer, что дает ему гибкость в управлении чтением (например, перечитать сообщения с начала).
Пример: Producer на Python (kafka-python)
from kafka import KafkaProducer
import json
# Создание Producer, подключающегося к кластеру Kafka
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
# Сериализация сообщений в JSON (в виде байтов)
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
# Отправка сообщения в топик 'user-events'
producer.send('user-events', {'user_id': 123, 'event': 'login'})
# Гарантируем, что все сообщения в буфере отправлены
producer.flush()
producer.close()
print("Сообщение успешно отправлено.")