6.2.3 · блок 6
MLOps-паттерны с Kafka
MLOps-паттерны с Kafka
Зачем это нужно
Kafka — не только «очередь сообщений». В ML-платформе она связывает продукт, data engineering и inference. Знание типовых паттернов помогает не изобретать архитектуру с нуля и избежать антипаттернов (например, синхронный scoring через Kafka request-reply без timeout policy).
Этот урок собирает рецепты, которые вы встретите в capstone и на стажировке.
Основные идеи
Pattern 1: Log enrichment → streaming features.
App → Kafka (raw events) → Stream processor → Online store / secondary topic
- Stream processor: Spark Structured Streaming, Flink, ksqlDB, custom Python.
- Output: агрегаты per entity (user, session, device).
- ML inference читает фичи из store, не из Kafka напрямую (latency, coupling).
Pattern 2: Change Data Capture (CDC) → feature refresh.
PostgreSQL → Debezium → Kafka → consumer → offline/online feature tables
- Каждое изменение строки CRM — event в Kafka.
- ML получает near-real-time updates без polling DB.
- Важно: ordering per primary key через partition key.
Pattern 3: Async inference (request–response через topics).
Client → topic inference.requests → Worker (MLServer) → topic inference.responses
- Плюсы: буферизация пиков, decoupling.
- Минусы: сложнее SLA, нужен correlation_id, TTL на requests.
- Не заменяет sync REST для interactive UI с p95 < 100 ms — гибрид: REST для online, Kafka для batch scoring миллионов записей.
Pattern 4: Event-driven retraining.
Monitor → topic ml.drift.detected → Argo Event / Workflow trigger → train pipeline
- Drift detector публикует event с метриками и dataset pointer.
- Orchestrator решает: auto-retrain vs human approval (MMS gate).
Pattern 5: Model lifecycle events (audit trail).
Topic ml.model.events: { type: "deployed", model_id, version, approver, ts }
- Подписчики: observability, compliance, downstream caches (invalidate old model).
Anti-patterns.
| Anti-pattern | Почему плохо |
|--------------|--------------|
| Giant messages (full images in Kafka) | Use object storage + reference in event |
| Kafka as primary DB | Retention limited; no complex queries |
| One topic for everything | No isolation, schema chaos |
| Consumer без dead-letter queue (DLQ) | Poison messages block pipeline |
| Sync request-reply без circuit breaker | Cascading latency |
Dead Letter Queue (DLQ). Сообщения, которые не парсятся или fail business validation, идут в topic.dlq для ручного разбора. ML pipeline не должен «зависнуть» на одном битом event.
Idempotency key. event_id или (source, offset) — consumer записывает processed ids (Redis/DB) или использует upsert в feature table.
Exactly-once feature write (практичный компромисс). At-least-once + idempotent upsert в online store часто проще, чем Kafka Transactions end-to-end.
Как это выглядит на практике
Fraud detection (упрощённо):
1. payments.transactions — все транзакции (Avro schema v3).
2. Flink job: window 5 min, features tx_count, avg_amount, geo_velocity.
3. Features → Redis + Kafka topic features.fraud.v1 (для audit).
4. Inference: REST /score читает Redis; fallback — last known from batch.
5. Drift: если null_rate фичи > 5% → event в ml.data.quality.
Capstone integration points:
- Strimzi topic definitions в GitOps repo.
- Consumer lag на Grafana (модуль 5).
- Feature schema в том же repo, что Feast definitions (6.4).
Design decision template (ADR preview):
- Context: нужен online scoring для 2k RPS.
- Option A: REST only, batch features T-1.
- Option B: Kafka streaming features + Redis.
- Decision: B для 3 realtime фич; остальные batch из Feast offline.
Что сделать после занятия
- [ ] Выберите учебный кейс и опишите один паттерн из урока (diagram + 5 строк текста).
- [ ] Для async inference выпишите 3 поля обязательного message envelope (correlation_id, ...).
- [ ] Назовите два anti-pattern'а, которые вы могли бы случайно спроектировать.