diff --git a/deployments/k8s/iot-kafka-consumer.yaml b/deployments/k8s/iot-kafka-consumer.yaml index 7cbc1df..c89365e 100644 --- a/deployments/k8s/iot-kafka-consumer.yaml +++ b/deployments/k8s/iot-kafka-consumer.yaml @@ -27,7 +27,7 @@ spec: containers: - name: kafka-consumer # Тот же образ что и оператор — все IoT бинари в одном образе. - image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.67 + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.68 imagePullPolicy: Always command: ["/iot-kafka-consumer"] env: diff --git a/deployments/k8s/iot-mqtt-bridge.yaml b/deployments/k8s/iot-mqtt-bridge.yaml index 5134fdb..97ba660 100644 --- a/deployments/k8s/iot-mqtt-bridge.yaml +++ b/deployments/k8s/iot-mqtt-bridge.yaml @@ -45,7 +45,7 @@ spec: - name: mqtt-bridge # Тот же образ что и оператор — оба бинаря в одном слое (manager + iot-mqtt-bridge). # При смене версии оператора — менять тег и здесь. - image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.67 + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.68 imagePullPolicy: Always command: ["/iot-mqtt-bridge"] env: diff --git a/deployments/k8s/kafka.yaml b/deployments/k8s/kafka.yaml index 3090aa9..daa02dd 100644 --- a/deployments/k8s/kafka.yaml +++ b/deployments/k8s/kafka.yaml @@ -1,4 +1,4 @@ -# Изменено: 2026-04-06 +# Изменено: 2026-04-06 (добавлен postStart hook для предсоздания топика iot.telemetry) # kafka.yaml — минимальный деплой Apache Kafka в KRaft mode (без Zookeeper). # Образ: apache/kafka (официальный, бесплатный). # Используется для IoT telemetry pipeline: mqtt-bridge → Kafka → iot-kafka-consumer → Postgres. diff --git a/go.mod b/go.mod index 6d50959..a07a8a7 100644 --- a/go.mod +++ b/go.mod @@ -3,11 +3,16 @@ module gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless go 1.25 require ( + github.com/eclipse/paho.mqtt.golang v1.5.1 + github.com/go-logr/logr v1.2.3 + github.com/google/uuid v1.6.0 github.com/gorilla/mux v1.8.1 github.com/lib/pq v1.11.2 github.com/minio/minio-go/v7 v7.0.99 github.com/onsi/ginkgo/v2 v2.6.0 github.com/onsi/gomega v1.24.1 + github.com/rabbitmq/amqp091-go v1.10.0 + github.com/segmentio/kafka-go v0.4.50 k8s.io/api v0.26.0 k8s.io/apimachinery v0.26.0 k8s.io/client-go v0.26.0 @@ -19,12 +24,10 @@ require ( github.com/cespare/xxhash/v2 v2.1.2 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/dustin/go-humanize v1.0.1 // indirect - github.com/eclipse/paho.mqtt.golang v1.5.1 // indirect github.com/emicklei/go-restful/v3 v3.9.0 // indirect github.com/evanphx/json-patch/v5 v5.6.0 // indirect github.com/fsnotify/fsnotify v1.6.0 // indirect github.com/go-ini/ini v1.67.0 // indirect - github.com/go-logr/logr v1.2.3 // indirect github.com/go-logr/zapr v1.2.3 // indirect github.com/go-openapi/jsonpointer v0.19.5 // indirect github.com/go-openapi/jsonreference v0.20.0 // indirect @@ -35,7 +38,6 @@ require ( github.com/google/gnostic v0.5.7-v3refs // indirect github.com/google/go-cmp v0.5.9 // indirect github.com/google/gofuzz v1.1.0 // indirect - github.com/google/uuid v1.6.0 // indirect github.com/gorilla/websocket v1.5.3 // indirect github.com/imdario/mergo v0.3.6 // indirect github.com/josharian/intern v1.0.0 // indirect @@ -57,9 +59,7 @@ require ( github.com/prometheus/client_model v0.3.0 // indirect github.com/prometheus/common v0.37.0 // indirect github.com/prometheus/procfs v0.8.0 // indirect - github.com/rabbitmq/amqp091-go v1.10.0 // indirect github.com/rs/xid v1.6.0 // indirect - github.com/segmentio/kafka-go v0.4.50 // indirect github.com/spf13/pflag v1.0.5 // indirect github.com/tinylib/msgp v1.6.1 // indirect go.uber.org/atomic v1.7.0 // indirect diff --git a/go.sum b/go.sum index b517b42..352788c 100644 --- a/go.sum +++ b/go.sum @@ -301,6 +301,12 @@ github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsT github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/tinylib/msgp v1.6.1 h1:ESRv8eL3u+DNHUoSAAQRE50Hm162zqAnBoGv9PzScPY= github.com/tinylib/msgp v1.6.1/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA= +github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.1.2 h1:FHX5I5B4i4hKRVRBCFRxq1iQRej7WO3hhBuJf+UUySY= +github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4= +github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8= +github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= github.com/yuin/goldmark v1.1.25/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.1.32/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= @@ -313,9 +319,8 @@ go.opencensus.io v0.22.4/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw= go.uber.org/atomic v1.7.0 h1:ADUqmZGgLDDfbSL9ZmPxKTybcoEYHgpYfELNoN+7hsw= go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc= go.uber.org/goleak v1.1.10/go.mod h1:8a7PlsEVH3e/a/GLqe5IIrQx6GzcnRmZEufDUTk4A7A= -go.uber.org/goleak v1.2.0 h1:xqgm/S+aQvhWFTtR0XK3Jvg7z8kGV8P4X14IzwN3Eqk= -go.uber.org/goleak v1.2.0/go.mod h1:XJYK+MuIchqpmGmUSAzotztawfKvYLUIgg7guXrwVUo= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/multierr v1.6.0 h1:y6IPFStTAIT5Ytl7/XYmHvzXQ7S3g/IeZW9hyZ5thw4= go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU= go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI= diff --git a/iot/cmd/kafka-consumer/main.go b/iot/cmd/kafka-consumer/main.go index 00cd5e6..8d7fd59 100644 --- a/iot/cmd/kafka-consumer/main.go +++ b/iot/cmd/kafka-consumer/main.go @@ -21,6 +21,7 @@ import ( "os/signal" "strings" "syscall" + "time" kafka "github.com/segmentio/kafka-go" @@ -71,6 +72,11 @@ func main() { 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, ","), @@ -148,3 +154,38 @@ func getEnvOrDefault(key, defaultVal string) string { } 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") +} diff --git a/iot/cmd/mqtt-bridge/main.go b/iot/cmd/mqtt-bridge/main.go index 68f68d1..a919819 100644 --- a/iot/cmd/mqtt-bridge/main.go +++ b/iot/cmd/mqtt-bridge/main.go @@ -82,9 +82,9 @@ func main() { // Kafka writer — асинхронный, с автоматическим созданием топика. kafkaWriter := &kafka.Writer{ - Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...), - Topic: iotTelemetryTopic, - Balancer: &kafka.LeastBytes{}, + Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...), + Topic: iotTelemetryTopic, + Balancer: &kafka.LeastBytes{}, // Позволяет продолжать работу при временной недоступности Kafka (буфер в памяти) Async: false, RequiredAcks: kafka.RequireOne,