mirror of
https://github.com/hikhvar/mqtt2prometheus.git
synced 2026-08-20 11:36:15 +00:00
As mentioned in https://github.com/hikhvar/mqtt2prometheus/issues/23, we do not use any logging framework at all. This was fine for getting the exporter startet. However, with inreasing load the logging must be configureable. This PR is a start to replace all instances of "log.Printf" with the zap logger. The current configuration parameters are the log level and the log format (console, json). We might expose the log configuration to the config file. But I think this is overkill for the current state of the exporter.
147 lines
4.0 KiB
Go
147 lines
4.0 KiB
Go
package metrics
|
|
|
|
import (
|
|
"fmt"
|
|
"go.uber.org/zap"
|
|
"strconv"
|
|
"time"
|
|
|
|
gojsonq "github.com/thedevsaddam/gojsonq/v2"
|
|
|
|
"github.com/eclipse/paho.mqtt.golang"
|
|
"github.com/hikhvar/mqtt2prometheus/pkg/config"
|
|
"github.com/prometheus/client_golang/prometheus"
|
|
)
|
|
|
|
type Ingest struct {
|
|
metricConfigs map[string][]config.MetricConfig
|
|
deviceIDRegex *config.Regexp
|
|
collector Collector
|
|
MessageMetric *prometheus.CounterVec
|
|
logger *zap.Logger
|
|
}
|
|
|
|
func NewIngest(collector Collector, metrics []config.MetricConfig, deviceIDRegex *config.Regexp) *Ingest {
|
|
cfgs := make(map[string][]config.MetricConfig)
|
|
for i := range metrics {
|
|
key := metrics[i].MQTTName
|
|
cfgs[key] = append(cfgs[key], metrics[i])
|
|
}
|
|
return &Ingest{
|
|
metricConfigs: cfgs,
|
|
deviceIDRegex: deviceIDRegex,
|
|
collector: collector,
|
|
MessageMetric: prometheus.NewCounterVec(
|
|
prometheus.CounterOpts{
|
|
Name: "received_messages",
|
|
Help: "received messages per topic and status",
|
|
}, []string{"status", "topic"},
|
|
),
|
|
logger: config.ProcessContext.Logger(),
|
|
}
|
|
}
|
|
|
|
// validMetric returns config matching the metric and deviceID
|
|
// Second return value indicates if config was found.
|
|
func (i *Ingest) validMetric(metric string, deviceID string) (config.MetricConfig, bool) {
|
|
for _, c := range i.metricConfigs[metric] {
|
|
if c.SensorNameFilter.Match(deviceID) {
|
|
return c, true
|
|
}
|
|
}
|
|
return config.MetricConfig{}, false
|
|
}
|
|
|
|
func (i *Ingest) store(topic string, payload []byte) error {
|
|
var mc MetricCollection
|
|
deviceID := i.deviceID(topic)
|
|
parsed := gojsonq.New().FromString(string(payload))
|
|
|
|
for path := range i.metricConfigs {
|
|
rawValue := parsed.Find(path)
|
|
parsed.Reset()
|
|
if rawValue == nil {
|
|
continue
|
|
}
|
|
m, err := i.parseMetric(path, deviceID, rawValue)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to parse valid metric value: %w", err)
|
|
}
|
|
m.Topic = topic
|
|
mc = append(mc, m)
|
|
}
|
|
|
|
i.collector.Observe(deviceID, mc)
|
|
return nil
|
|
}
|
|
|
|
func (i *Ingest) parseMetric(metricPath string, deviceID string, value interface{}) (Metric, error) {
|
|
cfg, cfgFound := i.validMetric(metricPath, deviceID)
|
|
if !cfgFound {
|
|
return Metric{}, nil
|
|
}
|
|
|
|
var metricValue float64
|
|
|
|
if boolValue, ok := value.(bool); ok {
|
|
if boolValue {
|
|
metricValue = 1
|
|
} else {
|
|
metricValue = 0
|
|
}
|
|
} else if strValue, ok := value.(string); ok {
|
|
|
|
// If string value mapping is defined, use that
|
|
if cfg.StringValueMapping != nil {
|
|
|
|
floatValue, ok := cfg.StringValueMapping.Map[strValue]
|
|
if ok {
|
|
metricValue = floatValue
|
|
} else if cfg.StringValueMapping.ErrorValue != nil {
|
|
metricValue = *cfg.StringValueMapping.ErrorValue
|
|
} else {
|
|
return Metric{}, fmt.Errorf("got unexpected string data '%s'", strValue)
|
|
}
|
|
|
|
} else {
|
|
|
|
// otherwise try to parse float
|
|
floatValue, err := strconv.ParseFloat(strValue, 64)
|
|
if err != nil {
|
|
return Metric{}, fmt.Errorf("got data with unexpectd type: %T ('%s') and failed to parse to float", value, value)
|
|
}
|
|
metricValue = floatValue
|
|
|
|
}
|
|
|
|
} else if floatValue, ok := value.(float64); ok {
|
|
metricValue = floatValue
|
|
} else {
|
|
return Metric{}, fmt.Errorf("got data with unexpectd type: %T ('%s')", value, value)
|
|
}
|
|
return Metric{
|
|
Description: cfg.PrometheusDescription(),
|
|
Value: metricValue,
|
|
ValueType: cfg.PrometheusValueType(),
|
|
IngestTime: time.Now(),
|
|
}, nil
|
|
}
|
|
|
|
func (i *Ingest) SetupSubscriptionHandler(errChan chan<- error) mqtt.MessageHandler {
|
|
return func(c mqtt.Client, m mqtt.Message) {
|
|
i.logger.Debug("Got message", zap.String("topic", m.Topic()), zap.String("payload", string(m.Payload())))
|
|
err := i.store(m.Topic(), m.Payload())
|
|
if err != nil {
|
|
errChan <- fmt.Errorf("could not store metrics '%s' on topic %s: %s", string(m.Payload()), m.Topic(), err.Error())
|
|
i.MessageMetric.WithLabelValues("storeError", m.Topic()).Inc()
|
|
return
|
|
}
|
|
i.MessageMetric.WithLabelValues("success", m.Topic()).Inc()
|
|
}
|
|
}
|
|
|
|
// deviceID uses the configured DeviceIDRegex to extract the device ID from the given mqtt topic path.
|
|
func (i *Ingest) deviceID(topic string) string {
|
|
return i.deviceIDRegex.GroupValue(topic, config.DeviceIDRegexGroup)
|
|
}
|