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

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 подключаются.

  • Internalmy-cluster-kafka-bootstrap:9092 внутри K8s.

  • External — NodePort, LoadBalancer, Ingress Route для producers вне кластера.

  • TLS — обязателен в prod; Strimzi генерирует certs через KafkaUser.

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.

  • Lag — сколько сообщений consumer не успел прочитать.

  • Растущий lag → feature pipeline отстаёт → stale features в inference.

  • Alert: lag > threshold 15 min → page on-call.

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:

  • TLS + SASL/SCRAM или mTLS через Strimzi certs.

  • ACL: producer crm-backend — write only; ml-feature-writer — read only.

  • No PII в topic name / log messages при debug.

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

  • Найдите в документации Strimzi минимальный пример Kafka CR и выпишите 3 поля, которые влияют на HA.

  • Опишите, как бы вы мониторили consumer lag для ML feature pipeline.

  • Составьте runbook из 4 шагов: «lag растёт 2 часа».

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