Добавлены метрики Prometheus и дашборды Grafana (этап 2.5)
Все три Go-сервиса (ingest, device-control, rule-engine) отдают /metrics в формате Prometheus: счётчики MQTT/RabbitMQ/ClickHouse операций, гистограммы длительности батч-флашей и диспатча правил, переходы устройств online/offline. Добавлены сервисы prometheus и grafana в docker-compose с провижининг конфигом (datasource + два готовых дашборда: "Ingest & Telemetry" и "Automation & Devices").
This commit is contained in:
@@ -5,16 +5,21 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/signal"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
|
||||
"git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/batch"
|
||||
chstore "git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/clickhouse"
|
||||
"git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/config"
|
||||
"git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/metrics"
|
||||
"git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/mqttclient"
|
||||
"git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/rabbitmq"
|
||||
"git.cactoz.su/cacto/home_automatization/services/ingest-service/internal/telemetry"
|
||||
@@ -59,7 +64,7 @@ func run(logger *slog.Logger) error {
|
||||
}
|
||||
defer publisher.Close()
|
||||
|
||||
batcher := batch.New(cfg.BatchMaxSize, cfg.BatchFlushInterval, cfg.BatchFlushTimeout, store.InsertBatch, logger)
|
||||
batcher := batch.New(cfg.BatchMaxSize, cfg.BatchFlushInterval, cfg.BatchFlushTimeout, instrumentedInsert(store), logger)
|
||||
batcher.Start()
|
||||
defer batcher.Stop()
|
||||
|
||||
@@ -77,6 +82,9 @@ func run(logger *slog.Logger) error {
|
||||
cancel()
|
||||
if err != nil {
|
||||
logger.Error("publish reading event failed", "device_id", reading.DeviceID, "error", err)
|
||||
metrics.RabbitMQPublishes.WithLabelValues("error").Inc()
|
||||
} else {
|
||||
metrics.RabbitMQPublishes.WithLabelValues("success").Inc()
|
||||
}
|
||||
}
|
||||
}()
|
||||
@@ -87,9 +95,12 @@ func run(logger *slog.Logger) error {
|
||||
Topic: cfg.MQTTTopic,
|
||||
QoS: 1,
|
||||
}, func(payload []byte, receivedAt time.Time) {
|
||||
metrics.MQTTMessagesReceived.Inc()
|
||||
|
||||
reading, err := telemetry.ParseReading(payload, receivedAt)
|
||||
if err != nil {
|
||||
logger.Warn("dropping invalid telemetry payload", "error", err)
|
||||
metrics.MQTTMessagesInvalid.Inc()
|
||||
return
|
||||
}
|
||||
|
||||
@@ -97,12 +108,24 @@ func run(logger *slog.Logger) error {
|
||||
case incoming <- reading:
|
||||
default:
|
||||
logger.Error("dropping reading: worker backlog full", "device_id", reading.DeviceID)
|
||||
metrics.MQTTMessagesDropped.Inc()
|
||||
}
|
||||
}, logger)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
metricsServer := &http.Server{
|
||||
Addr: fmt.Sprintf(":%d", cfg.MetricsPort),
|
||||
Handler: promhttp.Handler(),
|
||||
}
|
||||
go func() {
|
||||
logger.Info("metrics server started", "metrics_port", cfg.MetricsPort)
|
||||
if err := metricsServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
logger.Error("metrics server stopped", "error", err)
|
||||
}
|
||||
}()
|
||||
|
||||
logger.Info("ingest-service started", "mqtt_topic", cfg.MQTTTopic)
|
||||
|
||||
stop := make(chan os.Signal, 1)
|
||||
@@ -114,5 +137,30 @@ func run(logger *slog.Logger) error {
|
||||
close(incoming)
|
||||
workerWG.Wait()
|
||||
|
||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
if err := metricsServer.Shutdown(shutdownCtx); err != nil {
|
||||
logger.Error("metrics server shutdown failed", "error", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// instrumentedInsert wraps store.InsertBatch with duration/size/outcome
|
||||
// metrics, without teaching the batch or clickhouse packages about Prometheus.
|
||||
func instrumentedInsert(store *chstore.Store) batch.FlushFunc {
|
||||
return func(ctx context.Context, readings []telemetry.Reading) error {
|
||||
start := time.Now()
|
||||
err := store.InsertBatch(ctx, readings)
|
||||
|
||||
metrics.BatchFlushDuration.Observe(time.Since(start).Seconds())
|
||||
metrics.BatchSize.Observe(float64(len(readings)))
|
||||
if err != nil {
|
||||
metrics.BatchFlushes.WithLabelValues("error").Inc()
|
||||
} else {
|
||||
metrics.BatchFlushes.WithLabelValues("success").Inc()
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user