Files
IoT/legacy/doc/decisions/2026-04-12-replace-kafka-with-sqs.md
T

80 lines
3.6 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
> ⛔⛔⛔ ЛЕГАСИ (2026-08-16) — СТАРЫЙ IoT (k8s-деплой). НЕ ПРИНИМАТЬ ВО ВНИМАНИЕ. Актуальное: HISTORY/2026-08-16-session-log.md
# Решение: замена Kafka на shared-SQS
**Дата:** 2026-04-12
**Статус:** в работе
**Ветка:** `feature/replace-kafka-with-sqs`
---
## Контекст
IoT-сервис использует Kafka (библиотека `segmentio/kafka-go`) как промежуточную очередь между MQTT-bridge и Postgres consumer.
Kafka — тяжёлый компонент: требует отдельный деплоймент в кластере (`deployments/k8s/kafka.yaml`), ZooKeeper/KRaft, настройку топиков, партиций.
У нас есть собственный сервис **shared-SQS** — AWS SQS-совместимая очередь сообщений:
- Репа: https://gitea.services.ngcloud.ru/Nail/shared-SQS
- Endpoint: `https://qu.kube5s.ru`
- 17 SQS-операций, multi-tenant, Redis persistence, billing
- Работает через стандартные AWS SDK
RabbitMQ в коде IoT **не используется** (только в старых документах как план MVP).
---
## Решение
Заменить Kafka → shared-SQS во всём IoT pipeline.
---
## Что меняется
### 1. mqtt-bridge (`cmd/mqtt-bridge/main.go`)
- **Было:** `kafka.Writer` → `WriteMessages()` в топик `iot.telemetry`
- **Стало:** AWS SDK Go v2 → `sqs.SendMessage()` в очередь `iot-telemetry`
- Env: `KAFKA_BROKERS` → `SQS_ENDPOINT`, `SQS_ACCESS_KEY`, `SQS_SECRET_KEY`, `SQS_QUEUE_NAME`
### 2. kafka-consumer → sqs-consumer (`cmd/kafka-consumer/` → `cmd/sqs-consumer/`)
- **Было:** `kafka.NewReader` с consumer group, `FetchMessage()` + `CommitMessages()`
- **Стало:** polling loop: `sqs.ReceiveMessage(WaitTimeSeconds=20)` + `sqs.DeleteMessage()`
- Переименовать директорию и бинарник
### 3. admin stats handler (`internal/api/handler/iot_admin_stats_handler.go`)
- **Было:** Kafka lag (ListOffsets + OffsetFetch)
- **Стало:** `sqs.GetQueueAttributes(ApproximateNumberOfMessages, ApproximateNumberOfMessagesNotVisible)`
- Env: `KAFKA_BROKERS` → SQS credentials в handler
### 4. go.mod
- Убрать: `github.com/segmentio/kafka-go`
- Добавить: `github.com/aws/aws-sdk-go-v2`, `github.com/aws/aws-sdk-go-v2/service/sqs`
### 5. Deployments
- Удалить: `deployments/k8s/kafka.yaml`
- Обновить: `deployments/k8s/iot-mqtt-bridge.yaml` — новые env vars
- Обновить: `deployments/k8s/iot-kafka-consumer.yaml` → `iot-sqs-consumer.yaml`
### 6. Dockerfile / Makefile
- Переименовать бинарник `kafka-consumer` → `sqs-consumer`
---
## Плюсы
- Убираем Kafka из кластера (экономия ресурсов)
- Используем свой managed сервис (единая инфраструктура)
- AWS SDK — стандартная библиотека, код проще
- Billing и мониторинг из коробки в shared-SQS
## Риски
- SQS — pull-based (polling latency до 20s long poll vs Kafka push). Для IoT телеметрии приемлемо.
- At-least-once delivery — нужно учитывать idempotency (уже есть в текущем коде: INSERT ON CONFLICT)
---
## SQS endpoint
Пока используем публичный `https://qu.kube5s.ru`. Если есть внутрикластерный сервис — обновим.