Зачем это нужно
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 (упрощённо):
payments.transactions— все транзакции (Avro schema v3).Flink job: window 5 min, features
tx_count,avg_amount,geo_velocity.Features → Redis + Kafka topic
features.fraud.v1(для audit).Inference: REST
/scoreчитает Redis; fallback — last known from batch.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'а, которые вы могли бы случайно спроектировать.