docs: Kafka pipeline v0.1.68 — полная документация, race condition fix, тесты

This commit is contained in:
Naeel
2026-04-06 17:42:44 +03:00
parent 07ada8e362
commit 815b861417
2 changed files with 176 additions and 25 deletions
+139
View File
@@ -91,3 +91,142 @@
- `kafka-go` vs `confluent-kafka-go` — выбираем `segmentio/kafka-go` (pure Go, без CGO, совместим с alpine)
- KRaft mode в Helm bitnami — убедиться что включён (без Zookeeper)
- Topic `iot.telemetry` — создаётся автоматически при первой публикации (auto.create.topics.enable=true по умолчанию)
---
## Сессия (продолжение) — реализация Kafka pipeline
### Что было сделано
#### Ветка: `iot-kafka`
**1. Kafka StatefulSet (`deployments/k8s/kafka.yaml`)**
Установка через Helm bitnami провалилась — образ `bitnami/kafka:4.0.0` заблокирован (paywall с Aug 2025).
Переключились на официальный `apache/kafka:3.7.0` — бесплатный, полнофункциональный.
Написан кастомный `kafka.yaml`:
- KRaft mode (без Zookeeper) — node.id=1, roles=broker+controller
- ConfigMap монтируется в `/tmp/kafka-config` (не `/etc/kafka` — read-only в образе)
- `securityContext.fsGroup=1000` — kafka user (UID 1000) может писать в PVC
- PVC 1Gi на `vcd-disk-ext4` (local-path отказал: not enough disk space)
- Два Service: `kafka:9092` и headless `kafka-headless`
**2. bridge переписан (`iot/cmd/mqtt-bridge/main.go`)**
- Убран RabbitMQ (`amqp091-go`)
- Убрана прямая запись в Postgres через `iotpg`
- Добавлен Kafka writer (`segmentio/kafka-go`)
- Топик: `iot.telemetry`, ключ = namespace (партиционирование по тенанту)
- `Async: false, RequiredAcks: RequireOne` — синхронная запись, подтверждение от лидера
**3. kafka-consumer создан (`iot/cmd/kafka-consumer/main.go`)**
- Consumer group: `iot-pg-consumer`
- Читает из `iot.telemetry`, пишет в Postgres через `iotpg.Store`
- Offset коммитится ТОЛЬКО после успешной записи (at-least-once)
- Retry loop при недоступности Kafka
**4. Dockerfile обновлён**
- Добавлена сборка `iot-kafka-consumer` бинаря
- `COPY --from=builder /workspace/iot-kafka-consumer .`
- Итого в образе 3 бинаря: `manager`, `iot-mqtt-bridge`, `iot-kafka-consumer`
**5. Манифесты обновлены**
- `iot-mqtt-bridge.yaml`: убран `RABBITMQ_URL`, добавлен `KAFKA_BROKERS`
- `iot-kafka-consumer.yaml`: новый deployment
---
### Баги которые встретили и решили
#### Bug 1: дублирующий `package main`
`create_file` вставил `package main` дважды — в начале и перед `import`.
Фикс: `replace_string_in_file` удалил дубликат.
#### Bug 2: `kafka-go` помечен как `// indirect` в go.mod
gopls не видел пакет как доступный. Причина: зависимость добавлена без прямого импорта в момент добавления.
Фикс: `go mod tidy` убрал `// indirect`.
#### Bug 3: Race condition — consumer зависал при холодном старте
**Когда**: consumer стартовал одновременно с Kafka (первый деплой, топика нет).
**Что происходило**: consumer JOIN-ил group → Kafka auto-создавала топик в момент JOIN → kafka-go зависал на `FetchMessage` навсегда.
**Гипотеза №1**: postStart lifecycle hook на Kafka — создать топик сразу после старта брокера.
**Проблема с гипотезой**: `kafka-topics.sh --list` без таймаута зависает бесконечно → pod застрял в `PodInitializing`. Попытка с `nc` — `nc` не установлен в образе. Попытка с `request.timeout.ms` через properties — postStart возвращал exit code 1 → Kubernetes убивал контейнер → CrashLoopBackOff.
**Итоговое решение**: `ensureKafkaTopic()` в consumer — создаёт топик через `kafka.DialContext` + `conn.CreateTopics()` ДО создания Reader и JOIN группы. Retry 30 раз × 3 сек = 90 сек макс ожидания.
```go
// Порядок в consumer:
// 1. Connect IoT Postgres
// 2. ensureKafkaTopic() ← создаём топик, ждём брокер
// 3. kafka.NewReader() ← только теперь join group
// 4. FetchMessage() loop
```
**Почему это решение правильное**: race исключён на уровне приложения, не инфраструктуры. Даже если kafka.yaml не имеет никакого init — consumer сам дождётся Kafka и создаст топик.
#### Bug 4: CrashLoopBackOff после force delete pod-а
Force delete оставил `.lock` файл на PVC. Kafka падала с:
`Failed to acquire lock on file .lock in /var/kafka-data/logs`
Фикс: удалить StatefulSet + PVC (`kubectl delete statefulset kafka && kubectl delete pvc kafka-data-kafka-0`), пересоздать.
**Урок**: НИКОГДА не делать `kubectl delete pod --force` для stateful pod-ов. Только graceful (`kubectl delete pod`, подождать). Force delete = гарантированная поломка PVC.
---
### Результаты тестирования (v0.1.68)
| Тест | Условие | Результат |
|------|---------|-----------|
| Cold start | consumer стартует раньше Kafka | ✅ `ensureKafkaTopic` ретраится, дожидается |
| 5 рестартов consumer | Kafka работает | ✅ каждый раз `kafka topic ready` |
| MQTT → Pipeline | device2, 1 сообщение | ✅ offset=0 в Postgres |
| Рестарт Kafka | consumer живёт | ✅ ретраится с `ERROR fetch`, восстанавливается |
| 10 сообщений параллельно | 10 pod-ов mosquitto | ✅ offsets 2-11 все в Postgres |
**Что НЕ тестировалось:**
- Полный холодный старт с нуля (`kubectl apply -f` на чистый кластер)
- Consumer стартует одновременно с Kafka (оба новые) — race condition исправлен кодом, но на новом кластере не проверялся
---
### Текущее состояние кластера (2026-04-06 ~17:30 МСК)
```
sless-operator:v0.1.68 — Running
kafka-0 — Running (после удаления PVC и пересоздания)
iot-mqtt-bridge — Running, подключён к EMQX и Kafka
iot-kafka-consumer — Running, waiting for messages
iot-postgres — Running
```
Тенант: `sless-16367aacb67a4a01`, устройство `device2`.
В IoT Postgres: 12+ записей телеметрии (offsets 0-11).
---
### Что нужно сделать ещё
1. **Тест: полный холодный старт** — удалить kafka + consumer + PVC, применить всё одновременно, убедиться что race не вылезает
2. **Helm chart** — параметризовать `KAFKA_BROKERS`, `IOT_PG_DSN`, тег образа, StorageClass для `values-dev.yaml` / `values-prod.yaml`
3. **Managed Kafka/Postgres** — при переходе только менять `values-prod.yaml`
4. **Merge `iot-kafka` в `main`** — после тестов
---
### Архитектурные выводы сессии
**Будущая prod-архитектура (принято):**
- 3 кластера: IoT / Serverless / Infra-Control
- Managed Kafka + Managed Postgres (переключение через env vars, код не меняется)
- Helm chart для параметризации per-environment
**Текущий статус пути данных:**
```
IoT Device
→ MQTT PUBLISH
→ EMQX (sless namespace)
→ iot-mqtt-bridge (подписан на +/telemetry/+)
→ Kafka топик iot.telemetry (key=namespace)
→ iot-kafka-consumer (group iot-pg-consumer)
→ IoT Postgres (per-tenant schema через EnsureTenantDB)
→ GET /v1/{ns}/iot/telemetry (IoT Console)
```