6.2.1 · блок 6
Событийная модель и Apache Kafka
Событийная модель и Apache Kafka
Зачем это нужно
Многие 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).
1. Trigger retrain — счётчик drift events в topic → workflow стартует train.
2. Streaming features — агрегат «кликов за 5 мин» пишется в online store.
3. 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.