diff --git a/doc/progress.md b/doc/progress.md index 20ab09b..b895430 100644 --- a/doc/progress.md +++ b/doc/progress.md @@ -1,44 +1,56 @@ # Прогресс разработки -Последнее обновление: 2026-04-06 +Последнее обновление: 2026-04-06 17:30 МСК --- -## 2026-04-06 — Kafka интеграция (ветка iot-kafka, в процессе) +## 2026-04-06 — Kafka pipeline ЗАВЕРШЁН (v0.1.68, ветка iot-kafka) -### Цель -Заменить прямой INSERT в Postgres из bridge на Kafka pipeline: +### Итог +End-to-end IoT pipeline работает: ``` -MQTT → bridge → Kafka → iot-kafka-consumer → Postgres - → event-dispatcher → Functions (будущее) +MQTT Device → EMQX → iot-mqtt-bridge → Kafka → iot-kafka-consumer → IoT Postgres → GET /iot/telemetry ``` -### Обоснование -- RabbitMQ для IoT был подключён в bridge но бесполезен — никто не читал очередь -- Kafka даёт буферизацию, retention 7 дней, множество потребителей -- При переходе на managed Kafka в prod — только меняется KAFKA_BROKERS в Secret +### Что было сделано +- ✅ Kafka `apache/kafka:3.7.0` StatefulSet в KRaft mode (`deployments/k8s/kafka.yaml`) +- ✅ bridge переписан: убран RabbitMQ, добавлен Kafka producer +- ✅ `iot/cmd/kafka-consumer/main.go` — новый сервис, читает Kafka → пишет Postgres +- ✅ Dockerfile: 3 бинаря в одном образе (`manager`, `iot-mqtt-bridge`, `iot-kafka-consumer`) +- ✅ Race condition устранён: `ensureKafkaTopic()` создаёт топик до JOIN consumer group +- ✅ Тестирование: 5 рестартов consumer, рестарт Kafka, 10 сообщений параллельно +- ✅ Коммит `07ada8e`, образ `v0.1.68` в registry -### План -1. ✅ Документация + план -2. ✅ Ветка `iot-kafka` -3. ⏳ Helm: установить Kafka (bitnami, KRaft, 1 нод, PVC) в namespace `sless` -4. ⏳ bridge: убрать RabbitMQ, добавить Kafka producer (`segmentio/kafka-go`) -5. ⏳ Новый сервис `iot/cmd/kafka-consumer/main.go` -6. ⏳ Dockerfile + deployment манифесты -7. ⏳ Сборка v0.1.67, деплой, тест E2E +### Нерешённое +- ⚠️ Полный холодный старт (`kubectl apply -f` на чистый кластер) — НЕ ТЕСТИРОВАЛСЯ +- ⚠️ `rabbitmq` deployment в кластере — не используется IoT, можно убрать +- ⚠️ Helm chart — пока нет, нужен при переходе на managed Kafka/Postgres -### Что НЕ меняется -- EMQX, operator, REST API, IoT Console -- `iotpg` storage package -- ACL, auth, namespace-изоляция +### Версии +- Образ: `sless-operator:v0.1.68` +- Ветка: `iot-kafka` (коммит `07ada8e`) +- Kafka: `apache/kafka:3.7.0` (KRaft, 1 нод, PVC 1Gi на `vcd-disk-ext4`) + +### Ключевые уроки +1. **`kubectl delete pod --force` ломает PVC** у stateful pod-ов — оставляет `.lock` файл. Только graceful delete. +2. **postStart lifecycle hook** не подходит для "подождать пока сервис стартует" — нет `nc`, `kafka-topics.sh` зависает, exit code 1 убивает контейнер. +3. **Race condition kafka-go** при одновременном auto-create топика и join группы — решается предсозданием топика через admin API в consumer ДО создания Reader. +4. **`// indirect` в go.mod** = gopls не видит пакет. Фикс: `go mod tidy`. + +--- + +## 2026-04-06 (утро) — Kafka план + архитектурные решения + +### Принятые архитектурные решения +- 3 кластера в prod: IoT / Serverless / Infra-Control +- Managed Kafka + Managed Postgres (переключение через env vars) +- Helm chart нужен для параметризации per-environment +- `apache/kafka:3.7.0` вместо Bitnami (платный с Aug 2025 — НИКОГДА не упоминать) --- ## 2026-04-05 (вечер) — v0.1.66: UX-правки + деструктивный инцидент ---- - -## 2026-04-05 (вечер) — v0.1.66: UX-правки + деструктивный инцидент ### Изменения кода diff --git a/doc/thinking/2026-04-06.md b/doc/thinking/2026-04-06.md index 9abc325..30f4e1c 100644 --- a/doc/thinking/2026-04-06.md +++ b/doc/thinking/2026-04-06.md @@ -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) +```