doc: 2026-04-06 thinking log + kafka integration plan

This commit is contained in:
Naeel
2026-04-06 16:03:57 +03:00
parent 184f5ceb91
commit e46a8bb3e3
2 changed files with 128 additions and 1 deletions
+35 -1
View File
@@ -1,6 +1,40 @@
# Прогресс разработки
Последнее обновление: 2026-04-05
Последнее обновление: 2026-04-06
---
## 2026-04-06 — Kafka интеграция (ветка iot-kafka, в процессе)
### Цель
Заменить прямой INSERT в Postgres из bridge на Kafka pipeline:
```
MQTT → bridge → Kafka → iot-kafka-consumer → Postgres
→ event-dispatcher → Functions (будущее)
```
### Обоснование
- RabbitMQ для IoT был подключён в bridge но бесполезен — никто не читал очередь
- Kafka даёт буферизацию, retention 7 дней, множество потребителей
- При переходе на managed Kafka в prod — только меняется KAFKA_BROKERS в Secret
### План
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
### Что НЕ меняется
- EMQX, operator, REST API, IoT Console
- `iotpg` storage package
- ACL, auth, namespace-изоляция
---
## 2026-04-05 (вечер) — v0.1.66: UX-правки + деструктивный инцидент
---
+93
View File
@@ -0,0 +1,93 @@
# Thinking Log — 2026-04-06
## Агент: GitHub Copilot (Claude Sonnet 4.6)
---
## Архитектурные обсуждения перед началом Kafka
### Контекст
Пользователь обсуждал будущую prod-архитектуру IoT сервиса.
Никакого кода не менялось — чистое планирование.
### Итоги обсуждений
**Три отдельных кластера (принято):**
1. IoT кластер — EMQX, bridge, Kafka, iot-consumer, Postgres, REST API
2. Serverless кластер — operator, builder, event-dispatcher, Functions
3. Infra/Control кластер — Terraform для provisioning кластеров 1 и 2, DNS, TLS, auth, billing
Это классическая схема "control plane отдельно от data plane".
**Kafka — выбор подтверждён:**
- Сейчас: bridge → Postgres напрямую (синхронно, без буфера)
- Prod: bridge → Kafka → {consumer → Postgres, event-dispatcher → Functions}
- Dev/test: Kafka через Helm (bitnami, KRaft mode, 1 нод, PVC)
- Prod: managed Kafka (Confluent/Aiven) — только меняется KAFKA_BROKERS в Secret
**Postgres → managed облачный: легко**
- bridge и API используют DATABASE_URL из env
- Для переключения: только заменить Secret в кластере
- Код не трогается
**Состояние RabbitMQ для IoT (важное открытие):**
- Bridge сейчас пишет в RabbitMQ очередь `iot.{namespace}.telemetry`
- НО event-dispatcher эту очередь не читает — он настроен на serverless functions triggers
- То есть IoT-сообщения в RabbitMQ лежат мёртвым грузом — никто не читает
- Kafka заменяет RabbitMQ для IoT-части полностью
**Что проверяли в кластере:**
- 2026-04-05: только один активный тенант `sless-16367aacb67a4a01` (созданный после инцидента)
- Устройство `device2`, одно сообщение: `{"msg":"hello1dddd1777"}` от 14:34 UTC
- 2026-04-06: kubeconfig истёк → обновил → тот же один тенант, никто новый не входил
---
## План интеграции Kafka
### Анализ текущего bridge
Читал `iot/cmd/mqtt-bridge/main.go`. Текущая логика в `buildMQTTMessageHandler`:
1. Получает MQTT сообщение
2. Публикует в RabbitMQ (бесполезно — никто не читает)
3. Пишет напрямую в Postgres через iotpg.Store
С Kafka нужно:
1. Получает MQTT сообщение
2. Публикует в Kafka топик `iot.telemetry` (единый топик, namespace в payload)
3. Убрать прямой INSERT в Postgres из bridge
### Что создаётся заново
**`iot/cmd/kafka-consumer/main.go`** — новый сервис:
- Читает из Kafka топика `iot.telemetry`
- Пишет в Postgres (та же логика что сейчас в bridge)
- Consumer group: `iot-pg-consumer`
**Изменения в bridge:**
- Убрать RabbitMQ
- Добавить Kafka producer (библиотека `github.com/segmentio/kafka-go`)
- Env var: `KAFKA_BROKERS` вместо `RABBITMQ_URL`
**Новые env vars:**
- bridge: `KAFKA_BROKERS=kafka.sless.svc.cluster.local:9092`
- consumer: `KAFKA_BROKERS=...`, `IOT_PG_DSN=...`
### Что НЕ меняется
- EMQX, operator, REST API, IoT Console — не трогаются
- `iotpg` storage package — используется consumer-ом напрямую
- ACL, auth, namespace-изоляция — не меняются
### Порядок работы
1. Документация + коммит (сейчас)
2. Ветка `iot-kafka`
3. Helm: установить Kafka в namespace `sless`
4. Переписать bridge: убрать RabbitMQ, добавить Kafka producer
5. Создать `iot/cmd/kafka-consumer/main.go`
6. Обновить Dockerfile (добавить сборку consumer)
7. Обновить deployment манифесты
8. Сборка v0.1.67, деплой, тест
### Риски
- `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 по умолчанию)