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

3.6 KiB
Raw Blame History

ЛЕГАСИ (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-совместимая очередь сообщений:

RabbitMQ в коде IoT не используется (только в старых документах как план MVP).


Решение

Заменить Kafka → shared-SQS во всём IoT pipeline.


Что меняется

1. mqtt-bridge (cmd/mqtt-bridge/main.go)

  • Было: kafka.WriterWriteMessages() в топик iot.telemetry
  • Стало: AWS SDK Go v2 → sqs.SendMessage() в очередь iot-telemetry
  • Env: KAFKA_BROKERSSQS_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.yamliot-sqs-consumer.yaml

6. Dockerfile / Makefile

  • Переименовать бинарник kafka-consumersqs-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. Если есть внутрикластерный сервис — обновим.