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

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'а, которые вы могли бы случайно спроектировать.

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