- mqtt-bridge: kafka.Writer → SQS SendMessage (AWS SDK Go v2) - kafka-consumer → sqs-consumer: polling loop с ReceiveMessage/DeleteMessage - admin stats: Kafka lag → SQS GetQueueAttributes - Удалён kafka-consumer, добавлен sqs-consumer - go.mod: убран segmentio/kafka-go, добавлен aws-sdk-go-v2 - Dockerfile, Makefile: kafka-consumer → sqs-consumer - .gitignore: исправлен чтобы не игнорировать cmd/ директории - deployments: новый iot-sqs-consumer.yaml, обновлён mqtt-bridge - doc/decisions: задокументировано решение
78 lines
3.4 KiB
Markdown
78 lines
3.4 KiB
Markdown
# Решение: замена 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`. Если есть внутрикластерный сервис — обновим.
|