- 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: задокументировано решение
208 lines
6.8 KiB
Go
208 lines
6.8 KiB
Go
// Создано: 2026-04-12 (замена kafka-consumer → sqs-consumer)
|
|
// sqs-consumer/main.go — iot-sqs-consumer: читает IoT телеметрию из shared-SQS → пишет в Postgres.
|
|
//
|
|
// Роль в архитектуре:
|
|
// SQS очередь "iot-telemetry" → iot-sqs-consumer → IoT Postgres (per-tenant DB)
|
|
//
|
|
// Long polling (WaitTimeSeconds=20) — минимизирует пустые запросы к SQS.
|
|
// При временной недоступности Postgres — сообщения остаются в очереди (visibility timeout).
|
|
// DeleteMessage вызывается только после успешной записи в Postgres.
|
|
//
|
|
// Конфигурация через env vars:
|
|
// SQS_ENDPOINT — https://qu.kube5s.ru (или внутрикластерный endpoint)
|
|
// SQS_ACCESS_KEY — Access Key для shared-SQS tenant
|
|
// SQS_SECRET_KEY — Secret Key для shared-SQS tenant
|
|
// SQS_QUEUE_NAME — имя очереди (default: iot-telemetry)
|
|
// SQS_REGION — регион (default: us-east-1)
|
|
// IOT_PG_DSN — postgres://user:pass@host:5432/iotdb (master DSN для IoT Postgres)
|
|
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"log/slog"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/aws/aws-sdk-go-v2/aws"
|
|
"github.com/aws/aws-sdk-go-v2/credentials"
|
|
"github.com/aws/aws-sdk-go-v2/service/sqs"
|
|
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/storage/iotpg"
|
|
)
|
|
|
|
// sqsConsumerConfig — конфигурация из env vars.
|
|
type sqsConsumerConfig struct {
|
|
SQSEndpoint string
|
|
SQSAccessKey string
|
|
SQSSecretKey string
|
|
SQSQueueName string
|
|
SQSRegion string
|
|
}
|
|
|
|
// iotTelemetryMessage — envelope из SQS (идентичен bridge).
|
|
type iotTelemetryMessage struct {
|
|
Namespace string `json:"namespace"`
|
|
DeviceID string `json:"device_id"`
|
|
Topic string `json:"topic"`
|
|
Payload json.RawMessage `json:"payload"`
|
|
ReceivedAt string `json:"received_at"`
|
|
}
|
|
|
|
func main() {
|
|
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
|
|
|
|
cfg := loadConsumerConfig()
|
|
|
|
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
|
|
defer cancel()
|
|
|
|
log.Info("starting iot-sqs-consumer",
|
|
"sqs_endpoint", cfg.SQSEndpoint,
|
|
"sqs_queue", cfg.SQSQueueName,
|
|
)
|
|
|
|
// IoT Postgres — обязательный компонент для этого сервиса
|
|
iotStore, err := iotpg.NewFromEnv(log)
|
|
if err != nil || iotStore == nil {
|
|
log.Error("failed to connect to IoT Postgres — IOT_PG_DSN required", "err", err)
|
|
os.Exit(1)
|
|
}
|
|
defer iotStore.Close()
|
|
log.Info("connected to IoT Postgres")
|
|
|
|
// SQS клиент
|
|
sqsClient := sqs.New(sqs.Options{
|
|
Region: cfg.SQSRegion,
|
|
Credentials: credentials.NewStaticCredentialsProvider(
|
|
cfg.SQSAccessKey, cfg.SQSSecretKey, "",
|
|
),
|
|
BaseEndpoint: aws.String(cfg.SQSEndpoint),
|
|
})
|
|
|
|
// Получаем QueueUrl по имени один раз при старте
|
|
queueUrlOut, err := sqsClient.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{
|
|
QueueName: aws.String(cfg.SQSQueueName),
|
|
})
|
|
if err != nil {
|
|
log.Error("failed to get SQS queue URL — очередь должна существовать", "queue", cfg.SQSQueueName, "err", err)
|
|
os.Exit(1)
|
|
}
|
|
queueURL := *queueUrlOut.QueueUrl
|
|
log.Info("SQS queue resolved", "queue_url", queueURL)
|
|
|
|
log.Info("sqs consumer ready, starting polling loop...")
|
|
|
|
// Polling loop — long polling (WaitTimeSeconds=20) минимизирует пустые запросы.
|
|
// При ошибке SQS — backoff 5 сек и продолжаем.
|
|
for {
|
|
if ctx.Err() != nil {
|
|
break
|
|
}
|
|
|
|
resp, err := sqsClient.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
|
|
QueueUrl: aws.String(queueURL),
|
|
MaxNumberOfMessages: 10,
|
|
WaitTimeSeconds: 20, // long polling — SQS держит соединение до 20 сек
|
|
})
|
|
if err != nil {
|
|
if ctx.Err() != nil {
|
|
break // штатное завершение
|
|
}
|
|
log.Error("SQS ReceiveMessage failed", "err", err)
|
|
// Backoff при ошибках SQS — не спамим запросами
|
|
select {
|
|
case <-ctx.Done():
|
|
case <-time.After(5 * time.Second):
|
|
}
|
|
continue
|
|
}
|
|
|
|
for _, sqsMsg := range resp.Messages {
|
|
if sqsMsg.Body == nil {
|
|
continue
|
|
}
|
|
|
|
if err := processSQSTelemetry(ctx, *sqsMsg.Body, iotStore, log); err != nil {
|
|
log.Error("process telemetry message", "err", err, "message_id", derefStr(sqsMsg.MessageId))
|
|
// НЕ удаляем сообщение — оно вернётся в очередь после visibility timeout
|
|
continue
|
|
}
|
|
|
|
// Удаляем сообщение из SQS только после успешной записи в Postgres
|
|
_, err := sqsClient.DeleteMessage(ctx, &sqs.DeleteMessageInput{
|
|
QueueUrl: aws.String(queueURL),
|
|
ReceiptHandle: sqsMsg.ReceiptHandle,
|
|
})
|
|
if err != nil {
|
|
log.Error("SQS DeleteMessage failed", "err", err, "message_id", derefStr(sqsMsg.MessageId))
|
|
}
|
|
}
|
|
}
|
|
|
|
log.Info("shutting down iot-sqs-consumer")
|
|
}
|
|
|
|
// processSQSTelemetry десериализует envelope из SQS и записывает в Postgres.
|
|
func processSQSTelemetry(ctx context.Context, body string, store *iotpg.IoTPostgresStore, log *slog.Logger) error {
|
|
var envelope iotTelemetryMessage
|
|
if err := json.Unmarshal([]byte(body), &envelope); err != nil {
|
|
// Битое сообщение — логируем и пропускаем (не блокируем очередь)
|
|
log.Warn("failed to unmarshal telemetry envelope, skipping", "err", err, "raw", body)
|
|
return nil
|
|
}
|
|
|
|
// EnsureTenantDB идемпотентен — кэшируется после первого вызова
|
|
if err := store.EnsureTenantDB(ctx, envelope.Namespace); err != nil {
|
|
return err
|
|
}
|
|
|
|
if err := store.InsertTelemetry(ctx, envelope.Namespace, envelope.DeviceID, envelope.Payload); err != nil {
|
|
return err
|
|
}
|
|
|
|
log.Info("telemetry saved to Postgres",
|
|
"namespace", envelope.Namespace,
|
|
"device", envelope.DeviceID,
|
|
)
|
|
return nil
|
|
}
|
|
|
|
// loadConsumerConfig читает конфигурацию из env vars.
|
|
func loadConsumerConfig() sqsConsumerConfig {
|
|
required := func(key string) string {
|
|
v := os.Getenv(key)
|
|
if v == "" {
|
|
slog.Error("required env var not set", "key", key)
|
|
os.Exit(1)
|
|
}
|
|
return v
|
|
}
|
|
|
|
return sqsConsumerConfig{
|
|
SQSEndpoint: required("SQS_ENDPOINT"),
|
|
SQSAccessKey: required("SQS_ACCESS_KEY"),
|
|
SQSSecretKey: required("SQS_SECRET_KEY"),
|
|
SQSQueueName: getEnvOrDefault("SQS_QUEUE_NAME", "iot-telemetry"),
|
|
SQSRegion: getEnvOrDefault("SQS_REGION", "us-east-1"),
|
|
}
|
|
}
|
|
|
|
func getEnvOrDefault(key, defaultVal string) string {
|
|
if v := os.Getenv(key); v != "" {
|
|
return v
|
|
}
|
|
return defaultVal
|
|
}
|
|
|
|
// derefStr — безопасная разыменовка строкового указателя.
|
|
func derefStr(s *string) string {
|
|
if s == nil {
|
|
return ""
|
|
}
|
|
return *s
|
|
}
|