06 · Данные для ML6.1–6.4 · Потоки, батч и признаки6.2.1сложный

Событийная модель и 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_1h per 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.

Официальные материалы