From 387932ce1056b494033b69d07f0a677eea64c08b Mon Sep 17 00:00:00 2001 From: Naeel Date: Mon, 6 Apr 2026 18:26:03 +0300 Subject: [PATCH] =?UTF-8?q?fix:=20bridge=20Kafka=20write=20async=20(v0.1.6?= =?UTF-8?q?9)=20=E2=80=94=20MQTT=20callback=20=D0=BD=D0=B5=20=D0=B1=D0=BB?= =?UTF-8?q?=D0=BE=D0=BA=D0=B8=D1=80=D1=83=D0=B5=D1=82=D1=81=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- deployments/k8s/iot-kafka-consumer.yaml | 2 +- deployments/k8s/iot-mqtt-bridge.yaml | 2 +- doc/decisions/log.md | 34 ++++++++++++ doc/thinking/2026-04-06.md | 74 +++++++++++++++++++++++++ iot/cmd/mqtt-bridge/main.go | 28 +++++----- 5 files changed, 125 insertions(+), 15 deletions(-) diff --git a/deployments/k8s/iot-kafka-consumer.yaml b/deployments/k8s/iot-kafka-consumer.yaml index c89365e..7a8c42b 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.68 + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.69 imagePullPolicy: Always command: ["/iot-kafka-consumer"] env: diff --git a/deployments/k8s/iot-mqtt-bridge.yaml b/deployments/k8s/iot-mqtt-bridge.yaml index 97ba660..5495625 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.68 + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.69 imagePullPolicy: Always command: ["/iot-mqtt-bridge"] env: diff --git a/doc/decisions/log.md b/doc/decisions/log.md index d642a6f..6cc88f2 100644 --- a/doc/decisions/log.md +++ b/doc/decisions/log.md @@ -1255,3 +1255,37 @@ if err := h.K8s.Get(r.Context(), client.ObjectKey{...}, fn); err == nil { **Gap:** Для production нужен отдельный API-deployment с ≥2 replicas. +--- + +## 2026-04-06 — IoT bridge: Kafka write должен быть async (v0.1.69) + +### Контекст + +Load test (100 msg burst) показал потерю 73/100 сообщений. +Первоначально записал в "backlog". Пользователь указал: это не backlog — это +архитектурная ошибка. Между компонентами pipeline не должно быть синхронных зависимостей. + +### Решение + +`kafka.Writer{Async: true}` — единственно правильный вариант для MQTT callback. + +### Варианты которые рассматривались + +1. **`Async: true` в kafka.Writer** — выбрано. Минимальное изменение, kafka-go сам управляет буфером и горутиной записи. + +2. **Channel + отдельная горутина в handler** — избыточно. Дублирует то, что kafka-go уже делает внутри при Async=true. Лишний слой. + +3. **Увеличить keepalive timeout** — не решает проблему, только отодвигает симптом. + +### Почему `Async: true` безопасно + +- Ошибки доставки идут в `ErrorLogger` — логируются, не теряются бесследно +- При shutdown: `kafkaWriter.Close()` (defer) дожидается flush буфера перед выходом +- При недоступности Kafka: kafka-go внутри делает retry, сообщения в памяти-буфере + +### Принцип на будущее + +**Каждое звено pipeline должно принимать и отдавать сообщения немедленно.** +Любой blocking call внутри event handler — потенциальная точка потери данных. + + diff --git a/doc/thinking/2026-04-06.md b/doc/thinking/2026-04-06.md index fbf549d..bfd26eb 100644 --- a/doc/thinking/2026-04-06.md +++ b/doc/thinking/2026-04-06.md @@ -469,3 +469,77 @@ Consumer: iot-kafka-consumer-577f7ff88d-pkqd8, Running, 0 restarts Bridge: iot-mqtt-bridge-7dc87c46bc-tqjgz, Running, 0 restarts kafka-0: Running, 4 мин (перезапускался в TEST 7) ``` + +--- + +## Fix: v0.1.69 — Kafka write async (2026-04-06, после тестирования) +## Агент: GitHub Copilot (Claude Sonnet 4.6) + +### Проблема, выявленная тестом #3 + +При load test 100 сообщений выяснилось: **27/100 доставлено**. + +Первичная диагностика показала throughput ~1 msg/сек — я объяснил это +"bottleneck bridge" и записал в backlog. Но пользователь указал: это не backlog, +это архитектурная ошибка. **Между звеньями pipeline не должно быть ничего синхронного.** + +### Анализ root cause + +``` +MQTT callback (paho.mqtt.golang) вызывается синхронно в своём goroutine. +Если callback долго выполняется — следующие входящие MQTT сообщения накапливаются. +При Async=false: WriteMessages блокируется до получения ACK от Kafka (~1-10мс в норме, +но при burst + latency spike → сотни мс → EMQX keepalive timeout = disconnect). +``` + +Цепочка событий при burst: +1. 100 сообщений за <100мс влетают в EMQX +2. Bridge получает первое, вызывает WriteMessages (blocking ~1с) +3. Пока bridge заблокирован — EMQX keepalive не получает pingresp +4. После 30с (keepalive): EMQX разрывает соединение +5. Сообщения QoS 0, которые не были получены bridge — испаряются + +### Решение + +`kafka.Writer{Async: true}` — WriteMessages возвращается немедленно, Kafka batching +работает в фоновом goroutine внутри kafka-go. Ошибки доставки идут в `ErrorLogger`, +который логирует без блокировки MQTT loop. + +Почему **не** нужен отдельный channel/goroutine в handler: +kafka-go с `Async: true` уже внутри держит буфер и горутину записи. +Добавлять ещё один слой buffering — overengineering без причины. + +### Что изменено в коде (v0.1.69) + +**`iot/cmd/mqtt-bridge/main.go`:** +```go +// ДО (v0.1.68) — НЕПРАВИЛЬНО: +kafkaWriter := &kafka.Writer{ + Async: false, // блокирует MQTT callback до ACK Kafka +} +// в handler: +err = w.WriteMessages(ctx, ...) // блокировка ~1с/msg + +// ПОСЛЕ (v0.1.69) — ПРАВИЛЬНО: +kafkaWriter := &kafka.Writer{ + Async: true, // WriteMessages возвращается немедленно + ErrorLogger: kafka.LoggerFunc(func(msg string, args ...interface{}) { + log.Error("kafka async write error", ...) // ошибки не блокируют MQTT + }), +} +// в handler: +_ = w.WriteMessages(ctx, ...) // немедленный возврат, доставка в фоне +``` + +### Deployment manifests + +Оба yaml обновлены: `v0.1.68` → `v0.1.69`: +- `deployments/k8s/iot-mqtt-bridge.yaml` +- `deployments/k8s/iot-kafka-consumer.yaml` + +### Что ожидаем после фикса + +- MQTT callback завершается за <1мс (только marshal JSON + WriteMessages enqueue) +- Bridge не теряет keepalive с EMQX при burst +- Throughput: лимитируется сетью/Kafka, а не синхронным write (~тысячи msg/сек) +- Load test 100 сообщений: должны дойти все 100 diff --git a/iot/cmd/mqtt-bridge/main.go b/iot/cmd/mqtt-bridge/main.go index a919819..d5001c4 100644 --- a/iot/cmd/mqtt-bridge/main.go +++ b/iot/cmd/mqtt-bridge/main.go @@ -1,5 +1,5 @@ // Создано: 2026-04-04 -// Изменено: 2026-04-06 (заменён RabbitMQ на Kafka) +// Изменено: 2026-04-06 (fix: Kafka write async — MQTT callback не блокируется) // mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → Kafka. // // Роль в архитектуре: @@ -80,14 +80,18 @@ func main() { "kafka_brokers", cfg.KafkaBrokers, ) - // Kafka writer — асинхронный, с автоматическим созданием топика. + // Kafka writer — полностью асинхронный: WriteMessages возвращается немедленно, + // не блокируя MQTT callback. Kafka batching работает в фоне. + // Ошибки доставки логируются через ErrorLogger — не блокируют MQTT loop. kafkaWriter := &kafka.Writer{ - Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...), - Topic: iotTelemetryTopic, - Balancer: &kafka.LeastBytes{}, - // Позволяет продолжать работу при временной недоступности Kafka (буфер в памяти) - Async: false, + Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...), + Topic: iotTelemetryTopic, + Balancer: &kafka.LeastBytes{}, + Async: true, // MQTT callback не блокируется на ACK от Kafka RequiredAcks: kafka.RequireOne, + ErrorLogger: kafka.LoggerFunc(func(msg string, args ...interface{}) { + log.Error("kafka async write error", "detail", fmt.Sprintf(msg, args...)) + }), } defer kafkaWriter.Close() @@ -216,15 +220,13 @@ func buildMQTTMessageHandler(ctx context.Context, w *kafka.Writer, log *slog.Log } // Ключ = namespace — Kafka будет группировать сообщения одного тенанта - // на одну партицию (для упорядоченной обработки на consumer side) - err = w.WriteMessages(ctx, kafka.Message{ + // на одну партицию (для упорядоченной обработки на consumer side). + // WriteMessages с Async=true возвращается немедленно — не блокирует MQTT callback. + // Ошибки доставки идут в ErrorLogger выше. + _ = w.WriteMessages(ctx, kafka.Message{ Key: []byte(ns), Value: body, }) - if err != nil { - log.Error("publish to Kafka", "topic", iotTelemetryTopic, "namespace", ns, "err", err) - return - } log.Info("forwarded IoT telemetry to Kafka", "mqtt_topic", topic,