MLOps Path

6.2.2 · блок 6

Kafka на практике в Kubernetes

Kafka на практике в Kubernetes

Зачем это нужно

Теория topics и partitions бесполезна, если кластер Kafka не поднять, не подключить и не мониторить. В учебной и prod-среде курса Kafka обычно живёт в Kubernetes через Strimzi — так же declarative, как Deployment inference-сервиса.

Этот урок — мост от концепций к операционной работе: CRD, listeners, ACL, consumer lag, типовые поломки.

Основные идеи

Strimzi — Kafka Operator. Custom Resources:

| CRD | Назначение |

|-----|------------|

| `Kafka` | Кластер: brokers, версия, storage, listeners |

| `KafkaTopic` | Topic: partitions, replicas, retention |

| `KafkaUser` | ACL, TLS certificates для client |

| `KafkaConnect` | Connectors (S3 sink, JDBC source) |

| `KafkaMirrorMaker2` | Репликация между кластерами |

KRaft vs ZooKeeper. Современный Kafka (3.x+) может работать без ZooKeeper (KRaft mode). Strimzi поддерживает оба; для новых кластеров предпочтителен KRaft — меньше moving parts.

Listeners — как clients подключаются.

Storage. Kafka brokers пишут на persistent volumes (Longhorn, см. модуль 3.5). Потеря PV без replication = потеря данных partition. replication.factor >= 3, min.insync.replicas = 2 — типичный prod minimum.

KafkaTopic declarative:


apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaTopic
metadata:
  name: customer-events
  labels:
    strimzi.io/cluster: mdp-kafka
spec:
  partitions: 6
  replicas: 3
  config:
    retention.ms: 604800000   # 7 days
    cleanup.policy: delete

GitOps (Argo CD) хранит topic definitions в Git — не создаём topics руками в prod.

Consumer lag — главная метрика ML streaming.

Kafka UI / AKHQ / Redpanda Console — просмотр topics, messages, consumer groups. Для обучения — удобнее CLI:


# Список topics ( через kubectl exec в broker или kafka-bin )
kubectl -n kafka exec -it mdp-kafka-dual-role-0 -- \
  bin/kafka-topics.sh --bootstrap-server localhost:9092 --list

# Consumer group lag
kubectl -n kafka exec -it mdp-kafka-dual-role-0 -- \
  bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group feature-aggregator-v2

Типовые поломки.

| Симптом | Вероятная причина |

|---------|-------------------|

| Producer timeout | Broker down, wrong bootstrap, network policy |

| Rebalancing storm | Consumer slow / frequent restarts |

| Message too large | `max.message.bytes` vs payload size |

| Under-replicated partitions | Broker loss, disk full |

| Stale features in ML | Consumer lag, не Kafka itself |

Network policies. ML consumer Pod в namespace ml должен достучаться до bootstrap в kafka. Без NetworkPolicy или с неправильным — silent failure на connect.

Как это выглядит на практике

Учебный сценарий MDP:

1. Platform team деплоит Strimzi Kafka CR через Helm/Argo.

2. ML team MR: добавить KafkaTopic customer-events + KafkaUser ml-feature-writer.

3. Jenkins CI: integration test — producer/consumer в ephemeral namespace.

4. Feature job (Spark или Python faust/bytewax) — Deployment с env:

KAFKA_BOOTSTRAP=mdp-kafka-kafka-bootstrap.kafka.svc:9092

5. Grafana dashboard: consumer lag, bytes in/out per topic.

Local dev без полного кластера: Redpanda или Kafka в Docker Compose (модуль 3.1) — достаточно для прототипа consumer logic; перед prod — тест на Strimzi.

Security checklist для prod topic с PII events:

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

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

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