fix: bridge Kafka write async (v0.1.69) — MQTT callback не блокируется
This commit is contained in:
@@ -27,7 +27,7 @@ spec:
|
|||||||
containers:
|
containers:
|
||||||
- name: kafka-consumer
|
- name: kafka-consumer
|
||||||
# Тот же образ что и оператор — все IoT бинари в одном образе.
|
# Тот же образ что и оператор — все 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
|
imagePullPolicy: Always
|
||||||
command: ["/iot-kafka-consumer"]
|
command: ["/iot-kafka-consumer"]
|
||||||
env:
|
env:
|
||||||
|
|||||||
@@ -45,7 +45,7 @@ spec:
|
|||||||
- name: mqtt-bridge
|
- name: mqtt-bridge
|
||||||
# Тот же образ что и оператор — оба бинаря в одном слое (manager + iot-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
|
imagePullPolicy: Always
|
||||||
command: ["/iot-mqtt-bridge"]
|
command: ["/iot-mqtt-bridge"]
|
||||||
env:
|
env:
|
||||||
|
|||||||
@@ -1255,3 +1255,37 @@ if err := h.K8s.Get(r.Context(), client.ObjectKey{...}, fn); err == nil {
|
|||||||
|
|
||||||
**Gap:** Для production нужен отдельный API-deployment с ≥2 replicas.
|
**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 — потенциальная точка потери данных.
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -469,3 +469,77 @@ Consumer: iot-kafka-consumer-577f7ff88d-pkqd8, Running, 0 restarts
|
|||||||
Bridge: iot-mqtt-bridge-7dc87c46bc-tqjgz, Running, 0 restarts
|
Bridge: iot-mqtt-bridge-7dc87c46bc-tqjgz, Running, 0 restarts
|
||||||
kafka-0: Running, 4 мин (перезапускался в TEST 7)
|
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
|
||||||
|
|||||||
+15
-13
@@ -1,5 +1,5 @@
|
|||||||
// Создано: 2026-04-04
|
// Создано: 2026-04-04
|
||||||
// Изменено: 2026-04-06 (заменён RabbitMQ на Kafka)
|
// Изменено: 2026-04-06 (fix: Kafka write async — MQTT callback не блокируется)
|
||||||
// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → Kafka.
|
// mqtt-bridge/main.go — сервис-мост: MQTT (EMQX) → Kafka.
|
||||||
//
|
//
|
||||||
// Роль в архитектуре:
|
// Роль в архитектуре:
|
||||||
@@ -80,14 +80,18 @@ func main() {
|
|||||||
"kafka_brokers", cfg.KafkaBrokers,
|
"kafka_brokers", cfg.KafkaBrokers,
|
||||||
)
|
)
|
||||||
|
|
||||||
// Kafka writer — асинхронный, с автоматическим созданием топика.
|
// Kafka writer — полностью асинхронный: WriteMessages возвращается немедленно,
|
||||||
|
// не блокируя MQTT callback. Kafka batching работает в фоне.
|
||||||
|
// Ошибки доставки логируются через ErrorLogger — не блокируют MQTT loop.
|
||||||
kafkaWriter := &kafka.Writer{
|
kafkaWriter := &kafka.Writer{
|
||||||
Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...),
|
Addr: kafka.TCP(strings.Split(cfg.KafkaBrokers, ",")...),
|
||||||
Topic: iotTelemetryTopic,
|
Topic: iotTelemetryTopic,
|
||||||
Balancer: &kafka.LeastBytes{},
|
Balancer: &kafka.LeastBytes{},
|
||||||
// Позволяет продолжать работу при временной недоступности Kafka (буфер в памяти)
|
Async: true, // MQTT callback не блокируется на ACK от Kafka
|
||||||
Async: false,
|
|
||||||
RequiredAcks: kafka.RequireOne,
|
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()
|
defer kafkaWriter.Close()
|
||||||
|
|
||||||
@@ -216,15 +220,13 @@ func buildMQTTMessageHandler(ctx context.Context, w *kafka.Writer, log *slog.Log
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Ключ = namespace — Kafka будет группировать сообщения одного тенанта
|
// Ключ = namespace — Kafka будет группировать сообщения одного тенанта
|
||||||
// на одну партицию (для упорядоченной обработки на consumer side)
|
// на одну партицию (для упорядоченной обработки на consumer side).
|
||||||
err = w.WriteMessages(ctx, kafka.Message{
|
// WriteMessages с Async=true возвращается немедленно — не блокирует MQTT callback.
|
||||||
|
// Ошибки доставки идут в ErrorLogger выше.
|
||||||
|
_ = w.WriteMessages(ctx, kafka.Message{
|
||||||
Key: []byte(ns),
|
Key: []byte(ns),
|
||||||
Value: body,
|
Value: body,
|
||||||
})
|
})
|
||||||
if err != nil {
|
|
||||||
log.Error("publish to Kafka", "topic", iotTelemetryTopic, "namespace", ns, "err", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
log.Info("forwarded IoT telemetry to Kafka",
|
log.Info("forwarded IoT telemetry to Kafka",
|
||||||
"mqtt_topic", topic,
|
"mqtt_topic", topic,
|
||||||
|
|||||||
Reference in New Issue
Block a user