MLOps Path

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.

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:

Replay сценарий: DS нашёл баг в агрегации. Fix deploy → consumer reset offset на -7 days → пересчёт фич (осторожно с downstream idempotency).

Design review checklist:

Что сделать после занятия

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

Открыть интерактивную версию