- 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: задокументировано решение
3.4 KiB
3.4 KiB
Решение: замена 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. Если есть внутрикластерный сервис — обновим.