Зачем это нужно
Многие ML-сценарии живут в реальном времени: клик пользователя → обновить рекомендацию; транзакция → fraud score; лог приложения → anomaly detection. Batch-экспорт раз в сутки для таких задач слишком медленный. Событийная архитектура передаёт факты («что произошло») через шину сообщений; Apache Kafka — de facto стандарт такой шины в enterprise.
Понимание Kafka нужно MLOps-инженеру, чтобы проектировать streaming features, online inference triggers и интеграцию ML с продуктовыми сервисами.
Основные идеи
Event vs command.
Event — факт в прошлом:
OrderPlaced { order_id, user_id, amount, ts }. Неизменяем, append-only.Command — просьба что-то сделать:
SendEmail. В ML чаще работают с events как источником фич.
Kafka — распределённый commit log.
| Концепция | Аналогия | Зачем ML |
|---|---|---|
| Topic | Категория событий | user.clicks, payments.transactions |
| Partition | Упорядоченная очередь внутри topic | Параллелизм; key=user_id → все события пользователя в одной partition |
| Offset | Позиция в log | Replay для retrain / backfill |
| Producer | Пишет события | Backend, CDC connector |
| Consumer | Читает события | Spark Streaming, Flink, custom feature writer |
| Consumer Group | Группа consumers делят partitions | Horizontal scale без duplicate processing |
| Broker | Сервер Kafka | Кластер из 3+ broker для HA |
At-least-once vs exactly-once. Kafka по умолчанию at-least-once: сообщение может прийти дважды при retry. ML pipelines должны быть идемпотентными (dedup by event_id) или использовать exactly-once semantics (Kafka Transactions + Flink/Spark).
Retention. Сообщения хранятся N дней или до размера. Это не долговременное хранилище — для истории годами данные сливают в data lake (Kafka → S3/MinIO через Kafka Connect).
Schema Registry. События часто сериализуются Avro/Protobuf/JSON Schema. Registry хранит версии схем; consumer проверяет compatibility. Критично для ML: breaking schema change ломает feature pipeline.
Kafka vs message queue (RabbitMQ). Kafka оптимизирован под high throughput log и replay. RabbitMQ — routing, RPC patterns. Для ML feature streaming и event sourcing чаще Kafka.
Strimzi — оператор Kubernetes для Kafka (урок 6.2.2). Поднимает brokers, ZooKeeper/KRaft, listeners, TLS — declarative, как Deployment.
Event-driven ML patterns (preview).
Trigger retrain — счётчик drift events в topic → workflow стартует train.
Streaming features — агрегат «кликов за 5 мин» пишется в online store.
Async inference — request в topic, worker inference, response в другой topic.
Как это выглядит на практике
Topic customer.events, key = customer_id:
{ "event_id": "evt-8a2f", "event_type": "page_view", "customer_id": "c-9912", "page": "/pricing", "ts": "2026-01-15T14:32:01Z" }
Producer — веб-бэкенд через Kafka client или через sidecar.
Consumer group feature-aggregator-v2:
6 partitions → до 6 consumer instances.
Каждый считает rolling count
page_views_1hper customer.Результат → Redis (online features) + периодический snapshot в Parquet (offline).
Replay сценарий: DS нашёл баг в агрегации. Fix deploy → consumer reset offset на -7 days → пересчёт фич (осторожно с downstream idempotency).
Design review checklist:
Key выбран так, что нужный порядок per entity сохраняется?
Schema versioned в registry?
Retention достаточен для replay, но не бесконечен?
Monitoring: consumer lag, under-replicated partitions?
Что сделать после занятия
Нарисуйте topic с 3 partitions и 2 consumers в одной group — кто какую partition читает.
Придумайте event schema для учебного ML-кейса (3–5 полей + event_id).
Объясните, почему replay в Kafka полезен для ML, но опасен без idempotency.