MLOps Path

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

Pattern 2: Change Data Capture (CDC) → feature refresh.


PostgreSQL → Debezium → Kafka → consumer → offline/online feature tables

Pattern 3: Async inference (request–response через topics).


Client → topic inference.requests → Worker (MLServer) → topic inference.responses

Pattern 4: Event-driven retraining.


Monitor → topic ml.drift.detected → Argo Event / Workflow trigger → train pipeline

Pattern 5: Model lifecycle events (audit trail).


Topic ml.model.events: { type: "deployed", model_id, version, approver, ts }

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:

Design decision template (ADR preview):

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

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

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