feat: Kafka pipeline v0.1.67 (mqtt-bridge→kafka→consumer→postgres)

This commit is contained in:
Naeel
2026-04-06 16:41:31 +03:00
parent e46a8bb3e3
commit 63f834da2b
14 changed files with 885 additions and 135 deletions
+150
View File
@@ -0,0 +1,150 @@
// Создано: 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"
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")
// 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
}
+53 -121
View File
@@ -1,29 +1,26 @@
// Создано: 2026-04-04
// Изменено: 2026-04-05 (добавлен INSERT в IoT Postgres)
// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → RabbitMQ.
// Изменено: 2026-04-06 (заменён RabbitMQ на Kafka)
// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → Kafka.
//
// Роль в архитектуре:
// IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"] → RabbitMQ → event-dispatcher → function
// IoT Device → MQTT PUBLISH → EMQX → [mqtt-bridge подписан на "+/telemetry/+"]
// → Kafka топик "iot.telemetry"
// → iot-kafka-consumer → Postgres (история телеметрии)
// → event-dispatcher → Serverless Functions (триггеры)
//
// Логика:
// 1. Подключиться к EMQX как MQTT клиент (credentials из env)
// 2. Подписаться на топик "+/telemetry/+" (any namespace / telemetry / any device)
// 3. При получении сообщения:
// - Извлечь namespace из топика — первый сегмент до "/"
// - Опубликовать в RabbitMQ queue "iot.{namespace}.telemetry"
// - Payload передаётся as-is (JSON от устройства)
// 4. Переподключаться к RabbitMQ при разрыве (reconnect loop)
// 3. При получении сообщения — опубликовать в Kafka топик "iot.telemetry"
// 4. Payload оборачивается в envelope с метаданными (namespace, device_id, ts)
//
// Конфигурация через env vars:
// MQTT_BROKER_URL — tcp://emqx.sless.svc:1883
// MQTT_USERNAME — username для подключения bridge к EMQX
// MQTT_PASSWORD — пароль bridge клиента
// RABBITMQ_URL — amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/
// KAFKA_BROKERS — kafka.sless.svc.cluster.local:9092 (заменить на managed в prod)
//
// ВАЖНО: bridge клиент должен проходить EMQX auth — нужен IoTDevice "iot-bridge" в namespace "sless-bridge".
// Для MVP: выделить специальный namespace "sless-bridge" с устройством "bridge",
// и использовать его credentials для подключения bridge сервиса.
// Или: зарегистрировать bridge устройство через API и записать credentials в Secret.
// Для возврата к Postgres напрямую: см. git история, коммиты до 2026-04-06.
package main
@@ -39,9 +36,7 @@ import (
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
amqp "github.com/rabbitmq/amqp091-go"
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/iotpg"
kafka "github.com/segmentio/kafka-go"
)
// mqttBridgeConfig — конфигурация сервиса из env vars.
@@ -49,24 +44,28 @@ type mqttBridgeConfig struct {
MQTTBrokerURL string
MQTTUsername string
MQTTPassword string
RabbitMQURL string
KafkaBrokers string
}
// iotTelemetryMessage — структура сообщения публикуемого в RabbitMQ.
// Оборачивает MQTT payload в envelope с метаданными.
// iotTelemetryMessage — envelope сообщения публикуемого в Kafka.
// Потребители (iot-kafka-consumer, event-dispatcher) читают этот формат.
type iotTelemetryMessage struct {
// Namespace — k8s namespace пользователя (из MQTT topic)
// Namespace — k8s namespace тенанта (из MQTT topic, первый сегмент)
Namespace string `json:"namespace"`
// DeviceID — идентификатор устройства (из MQTT topic, последний сегмент)
// DeviceID — идентификатор устройства (из MQTT topic, третий сегмент)
DeviceID string `json:"device_id"`
// Topic — оригинальный MQTT topic
Topic string `json:"topic"`
// Payload — данные от устройства (JSON передаётся as-is / строка если не JSON)
// Payload — данные от устройства (JSON as-is, или строка если не JSON)
Payload json.RawMessage `json:"payload"`
// ReceivedAt — время получения сообщения мостом (UTC)
// ReceivedAt — время получения сообщения мостом (UTC, RFC3339)
ReceivedAt string `json:"received_at"`
}
// iotTelemetryTopic — Kafka топик для IoT телеметрии.
// Все устройства всех тенантов пишут в один топик, изоляция — по полю Namespace в payload.
const iotTelemetryTopic = "iot.telemetry"
func main() {
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
@@ -78,34 +77,19 @@ func main() {
log.Info("starting iot-mqtt-bridge",
"mqtt_broker", cfg.MQTTBrokerURL,
"mqtt_username", cfg.MQTTUsername,
"kafka_brokers", cfg.KafkaBrokers,
)
// RabbitMQ connection с reconnect loop
rabbitConn, err := connectRabbitMQWithRetry(ctx, cfg.RabbitMQURL, log)
if err != nil {
log.Error("failed to connect to RabbitMQ", "err", err)
os.Exit(1)
}
defer rabbitConn.Close()
rabbitCh, err := rabbitConn.Channel()
if err != nil {
log.Error("failed to open RabbitMQ channel", "err", err)
os.Exit(1)
}
defer rabbitCh.Close()
// IoT Postgres — сохранение телеметрии (per-tenant DB).
// Опционально: если IOT_PG_DSN не задан — продолжаем работать без Postgres (только RabbitMQ)
iotPGStore, err := iotpg.NewFromEnv(log)
if err != nil {
log.Error("failed to connect to IoT Postgres", "err", err)
os.Exit(1)
}
if iotPGStore != nil {
defer iotPGStore.Close()
log.Info("connected to IoT Postgres for telemetry storage")
// Kafka writer — асинхронный, с автоматическим созданием топика.
kafkaWriter := &kafka.Writer{
Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...),
Topic: iotTelemetryTopic,
Balancer: &kafka.LeastBytes{},
// Позволяет продолжать работу при временной недоступности Kafka (буфер в памяти)
Async: false,
RequiredAcks: kafka.RequireOne,
}
defer kafkaWriter.Close()
// Создаём MQTT клиент
mqttClient, err := connectMQTT(cfg, log)
@@ -116,8 +100,7 @@ func main() {
defer mqttClient.Disconnect(500)
// Функция-обработчик MQTT сообщений
// Вызывается в goroutine paho при каждом сообщении
messageHandler := buildMQTTMessageHandler(ctx, rabbitCh, iotPGStore, log)
messageHandler := buildMQTTMessageHandler(ctx, kafkaWriter, log)
// Подписываемся на все telemetry топики всех namespace
// "+/telemetry/+" = {любой namespace}/telemetry/{любой deviceId}
@@ -135,7 +118,6 @@ func main() {
}
// loadBridgeConfig читает конфигурацию из env vars.
// Завершает процесс если обязательные переменные отсутствуют.
func loadBridgeConfig() mqttBridgeConfig {
required := func(key string) string {
v := os.Getenv(key)
@@ -150,7 +132,7 @@ func loadBridgeConfig() mqttBridgeConfig {
MQTTBrokerURL: getEnvOrDefault("MQTT_BROKER_URL", "tcp://emqx.sless.svc:1883"),
MQTTUsername: required("MQTT_USERNAME"),
MQTTPassword: required("MQTT_PASSWORD"),
RabbitMQURL: required("RABBITMQ_URL"),
KafkaBrokers: getEnvOrDefault("KAFKA_BROKERS", "kafka.sless.svc.cluster.local:9092"),
}
}
@@ -162,7 +144,6 @@ func getEnvOrDefault(key, defaultVal string) string {
}
// connectMQTT устанавливает подключение к EMQX брокеру.
// AutoReconnect=true — paho сам переподключается при разрыве.
func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) {
opts := mqtt.NewClientOptions()
opts.AddBroker(cfg.MQTTBrokerURL)
@@ -173,7 +154,7 @@ func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) {
opts.SetConnectRetry(true)
opts.SetConnectRetryInterval(5 * time.Second)
opts.SetKeepAlive(30 * time.Second)
opts.SetCleanSession(false) // сохраняем подписки при реконнекте
opts.SetCleanSession(false)
opts.SetConnectionLostHandler(func(_ mqtt.Client, err error) {
log.Warn("MQTT connection lost, reconnecting...", "err", err)
@@ -187,7 +168,6 @@ func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) {
client := mqtt.NewClient(opts)
token := client.Connect()
// Ждём максимум 30 секунд
if !token.WaitTimeout(30 * time.Second) {
return nil, fmt.Errorf("MQTT connect timeout")
}
@@ -197,40 +177,15 @@ func connectMQTT(cfg mqttBridgeConfig, log *slog.Logger) (mqtt.Client, error) {
return client, nil
}
// connectRabbitMQWithRetry подключается к RabbitMQ с повторными попытками.
// Retry нужен потому что RabbitMQ может стартовать позже bridge сервиса.
func connectRabbitMQWithRetry(ctx context.Context, url string, log *slog.Logger) (*amqp.Connection, error) {
const maxAttempts = 10
for attempt := 1; attempt <= maxAttempts; attempt++ {
conn, err := amqp.Dial(url)
if err == nil {
log.Info("connected to RabbitMQ", "attempt", attempt)
return conn, nil
}
log.Warn("RabbitMQ connection failed, retrying...", "attempt", attempt, "err", err)
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(5 * time.Second):
}
}
return nil, fmt.Errorf("exhausted %d RabbitMQ connection attempts", maxAttempts)
}
// buildMQTTMessageHandler возвращает функцию-обработчик MQTT сообщений.
// Замыкание над rabbitCh (RabbitMQ channel), iotStore (может быть nil) и logger.
// Порядок действий при получении сообщения:
// 1. INSERT в IoT Postgres (tenant DB) — если iotStore != nil
// 2. Publish в RabbitMQ — всегда (для event-dispatcher → function triggers)
//
// Ошибка INSERT не блокирует RabbitMQ publish — разные failure domain.
func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotStore *iotpg.IoTPostgresStore, log *slog.Logger) mqtt.MessageHandler {
// buildMQTTMessageHandler возвращает обработчик MQTT сообщений.
// При получении сообщения — публикует envelope в Kafka топик "iot.telemetry".
// Ключ сообщения Kafka = namespace, для партиционирования по тенанту.
func buildMQTTMessageHandler(ctx context.Context, w *kafka.Writer, log *slog.Logger) mqtt.MessageHandler {
return func(_ mqtt.Client, msg mqtt.Message) {
topic := msg.Topic()
payload := msg.Payload()
// Топик: "{namespace}/telemetry/{deviceId}"
// Извлекаем namespace (первый сегмент) и deviceId (третий сегмент)
parts := strings.SplitN(topic, "/", 3)
if len(parts) != 3 {
log.Warn("unexpected MQTT topic format, skipping", "topic", topic)
@@ -239,28 +194,13 @@ func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotSto
ns := parts[0]
deviceID := parts[2]
// Нормализуем payload: если это не JSON — оборачиваем в строку
// Нормализуем payload: если не JSON — оборачиваем в строку
rawPayload := json.RawMessage(payload)
if !json.Valid(payload) {
quotedBytes, _ := json.Marshal(string(payload))
rawPayload = json.RawMessage(quotedBytes)
}
// ШАГ 1: INSERT в IoT Postgres — сохраняем телеметрию в per-tenant DB
// EnsureTenantDB идемпотентен: кэшируется после первого вызова
if iotStore != nil {
if err := iotStore.EnsureTenantDB(ctx, ns); err != nil {
log.Error("ensure tenant DB", "namespace", ns, "err", err)
// НЕ возвращаемся — продолжаем RabbitMQ publish
} else if err := iotStore.InsertTelemetry(ctx, ns, deviceID, rawPayload); err != nil {
log.Error("insert telemetry", "topic", topic, "err", err)
// НЕ возвращаемся — RabbitMQ не должен зависеть от Postgres
} else {
log.Debug("telemetry saved to Postgres", "namespace", ns, "device", deviceID)
}
}
// ШАГ 2: Publish в RabbitMQ (для event-dispatcher → function triggers)
envelope := iotTelemetryMessage{
Namespace: ns,
DeviceID: deviceID,
@@ -275,30 +215,22 @@ func buildMQTTMessageHandler(ctx context.Context, rabbitCh *amqp.Channel, iotSto
return
}
// Queue name: "iot.{namespace}.telemetry"
// Declare-on-publish: если queue не существует — создаём
queueName := fmt.Sprintf("iot.%s.telemetry", ns)
if _, err := rabbitCh.QueueDeclare(queueName, true, false, false, false, nil); err != nil {
log.Error("declare RabbitMQ queue", "queue", queueName, "err", err)
return
}
err = rabbitCh.Publish(
"", // exchange — default exchange
queueName, // routing key = queue name для default exchange
false, // mandatory
false, // immediate
amqp.Publishing{
ContentType: "application/json",
Body: body,
DeliveryMode: amqp.Persistent, // сохранять при рестарте RabbitMQ
},
)
// Ключ = namespace — Kafka будет группировать сообщения одного тенанта
// на одну партицию (для упорядоченной обработки на consumer side)
err = w.WriteMessages(ctx, kafka.Message{
Key: []byte(ns),
Value: body,
})
if err != nil {
log.Error("publish to RabbitMQ", "queue", queueName, "err", err)
log.Error("publish to Kafka", "topic", iotTelemetryTopic, "namespace", ns, "err", err)
return
}
log.Info("forwarded IoT telemetry", "topic", topic, "namespace", ns, "device", deviceID, "queue", queueName)
log.Info("forwarded IoT telemetry to Kafka",
"mqtt_topic", topic,
"namespace", ns,
"device", deviceID,
"kafka_topic", iotTelemetryTopic,
)
}
}