192 lines
7.0 KiB
Go
192 lines
7.0 KiB
Go
// Создано: 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")
|
||
}
|