// Создано: 2026-04-06 // kafka-consumer/main.go — iot-kafka-consumer: читает IoT телеметрию из Kafka → пишет в Postgres. // // Роль в архитектуре: // Kafka топик "iot.telemetry" → iot-kafka-consumer → IoT Postgres (per-tenant DB) // // Consumer group "iot-pg-consumer" — позволяет запускать несколько реплик без дублирования. // При временной недоступности Postgres — Kafka хранит сообщения (retention 7 дней). // // Конфигурация через env vars: // KAFKA_BROKERS — kafka.sless.svc.cluster.local:9092 (или managed Kafka в prod) // IOT_PG_DSN — postgres://user:pass@host:5432/iotdb (master DSN для IoT Postgres) package main import ( "context" "encoding/json" "log/slog" "os" "os/signal" "strings" "syscall" "time" kafka "github.com/segmentio/kafka-go" "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/iotpg" ) // kafkaConsumerConfig — конфигурация из env vars. type kafkaConsumerConfig struct { KafkaBrokers string } // iotTelemetryMessage — envelope из Kafka (идентичен 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"` } // iotTelemetryTopic — Kafka топик (должен совпадать с bridge). const iotTelemetryTopic = "iot.telemetry" // iotConsumerGroup — идентификатор consumer group. // При нескольких репликах Kafka распределяет партиции между ними. const iotConsumerGroup = "iot-pg-consumer" 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-kafka-consumer", "kafka_brokers", cfg.KafkaBrokers, "topic", iotTelemetryTopic, "group", iotConsumerGroup, ) // 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") // Предсоздаём топик ДО присоединения к consumer group. // Это устраняет race condition в kafka-go: если consumer joinит группу в момент // когда топик auto-создаётся — kafka-go зависает. Явное создание до Join это исключает. ensureKafkaTopic(ctx, cfg.KafkaBrokers, log) // Kafka reader с consumer group — автоматически коммитит offsets после обработки reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: strings.Split(cfg.KafkaBrokers, ","), Topic: iotTelemetryTopic, GroupID: iotConsumerGroup, MinBytes: 1, MaxBytes: 1 << 20, // 1MB }) defer reader.Close() log.Info("kafka reader ready, waiting for messages...") for { // FetchMessage — блокирует до следующего сообщения kafkaMsg, err := reader.FetchMessage(ctx) if err != nil { if ctx.Err() != nil { break // штатное завершение } log.Error("fetch from Kafka", "err", err) continue } if err := processKafkaTelemetry(ctx, kafkaMsg, iotStore, log); err != nil { log.Error("process telemetry message", "err", err) // НЕ коммитим offset — сообщение будет перечитано при следующем старте continue } // Коммитим offset только после успешной записи в Postgres if err := reader.CommitMessages(ctx, kafkaMsg); err != nil { log.Error("commit Kafka offset", "err", err) } } log.Info("shutting down iot-kafka-consumer") } // processKafkaTelemetry десериализует сообщение из Kafka и записывает в Postgres. func processKafkaTelemetry(ctx context.Context, msg kafka.Message, store *iotpg.IoTPostgresStore, log *slog.Logger) error { var envelope iotTelemetryMessage if err := json.Unmarshal(msg.Value, &envelope); err != nil { // Битое сообщение — логируем и пропускаем (не блокируем очередь) log.Warn("failed to unmarshal telemetry envelope, skipping", "err", err, "raw", string(msg.Value)) 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, "kafka_offset", msg.Offset, ) return nil } // loadConsumerConfig читает конфигурацию из env vars. func loadConsumerConfig() kafkaConsumerConfig { return kafkaConsumerConfig{ KafkaBrokers: getEnvOrDefault("KAFKA_BROKERS", "kafka.sless.svc.cluster.local:9092"), } } func getEnvOrDefault(key, defaultVal string) string { if v := os.Getenv(key); v != "" { return v } return defaultVal } // ensureKafkaTopic создаёт топик iot.telemetry если не существует. // Вызывается ДО создания Reader и Join consumer group — исключает race condition // в kafka-go при одновременном auto-create топика и join группы. // Ретраится пока Kafka не ответит (брокер может ещё стартовать). func ensureKafkaTopic(ctx context.Context, brokers string, log *slog.Logger) { brokerList := strings.Split(brokers, ",") for attempt := 1; attempt <= 30; attempt++ { conn, err := kafka.DialContext(ctx, "tcp", brokerList[0]) if err != nil { log.Warn("kafka not reachable yet, retrying...", "attempt", attempt, "err", err) select { case <-ctx.Done(): return case <-time.After(3 * time.Second): continue } } defer conn.Close() // Создаём топик идемпотентно — ошибка TopicAlreadyExists игнорируется err = conn.CreateTopics(kafka.TopicConfig{ Topic: iotTelemetryTopic, NumPartitions: 1, ReplicationFactor: 1, }) if err != nil && err != kafka.TopicAlreadyExists { log.Warn("failed to create kafka topic, auto.create.topics.enable will handle it", "err", err) } else { log.Info("kafka topic ready", "topic", iotTelemetryTopic) } return } log.Warn("kafka did not respond after 30 attempts, proceeding without pre-creation") }