From 0f07cd612c42c309c34aff404e349dbcf387afdd Mon Sep 17 00:00:00 2001 From: andig Date: Mon, 23 Jan 2023 18:48:34 +0100 Subject: [PATCH] InfluxDB: write slice of structs (#5873) --- server/influxdb.go | 93 +++++++++++++++++++++++++--------------------- 1 file changed, 50 insertions(+), 43 deletions(-) diff --git a/server/influxdb.go b/server/influxdb.go index 938ed6482..ba5f69e52 100644 --- a/server/influxdb.go +++ b/server/influxdb.go @@ -2,12 +2,16 @@ package server import ( "fmt" + "reflect" + "strconv" + "strings" "sync" "time" "github.com/evcc-io/evcc/core/loadpoint" "github.com/evcc-io/evcc/util" influxdb2 "github.com/influxdata/influxdb-client-go/v2" + "github.com/influxdata/influxdb-client-go/v2/api" influxlog "github.com/influxdata/influxdb-client-go/v2/log" ) @@ -54,22 +58,10 @@ func NewInfluxClient(url, token, org, user, password, database string) *Influx { } } -// supportedType checks if type can be written as influx value -func (m *Influx) supportedType(p util.Param) bool { - if p.Val == nil { - return true - } - - switch val := p.Val.(type) { - case int, int64, float64: - return true - case [3]float64: - return true - case []float64: - return len(val) == 3 - default: - return false - } +// writePoint asynchronously writes a point to influx +func (m *Influx) writePoint(writer api.WriteAPI, key string, fields map[string]any, tags map[string]string) { + m.log.TRACE.Printf("write %s=%v (%v)", key, fields, tags) + writer.WritePoint(influxdb2.NewPoint(key, tags, fields, time.Now())) } // Run Influx publisher @@ -90,49 +82,64 @@ func (m *Influx) Run(loadPoints []loadpoint.API, in <-chan util.Param) { // add points to batch for async writing for param := range in { // vehicle name - if param.Loadpoint != nil { - if name, ok := param.Val.(string); ok && param.Key == "vehicleTitle" { - vehicles[*param.Loadpoint] = name + if param.Loadpoint != nil && param.Key == "vehicleTitle" { + if vehicle, ok := param.Val.(string); ok { + vehicles[*param.Loadpoint] = vehicle continue } } - if !m.supportedType(param) { - continue - } + fields := make(map[string]any) - tags := map[string]string{} + tags := make(map[string]string) if param.Loadpoint != nil { tags["loadpoint"] = loadPoints[*param.Loadpoint].Name() tags["vehicle"] = vehicles[*param.Loadpoint] } - fields := map[string]interface{}{} + switch val := param.Val.(type) { + case int, int64, float64: + fields["value"] = param.Val - // array to slice - val := param.Val - if v, ok := val.([3]float64); ok { - val = v[:] - } - - // add slice as phase values - if phases, ok := val.([]float64); ok { - var total float64 - for i, v := range phases { - total += v + case [3]float64: + // add array as phase values + for i, v := range val { fields[fmt.Sprintf("l%d", i+1)] = v } - // add total as "value" - val = total + default: + // allow writing nil values + if param.Val == nil { + break + } + + // slice of structs + if typ := reflect.TypeOf(param.Val); typ.Kind() == reflect.Slice && typ.Elem().Kind() == reflect.Struct { + val := reflect.ValueOf(param.Val) + + // loop slice + for i := 0; i < val.Len(); i++ { + val := val.Index(i) + typ := val.Type() + + // loop struct + for j := 0; j < typ.NumField(); j++ { + n := typ.Field(j).Name + v := val.Field(j).Interface() + + key := param.Key + strings.ToUpper(n[:1]) + n[1:] + fields["value"] = v + tags["id"] = strconv.Itoa(i + 1) + + m.writePoint(writer, key, fields, tags) + } + } + } + + continue } - fields["value"] = val - - // write asynchronously - m.log.TRACE.Printf("write %s=%v (%v)", param.Key, param.Val, tags) - p := influxdb2.NewPoint(param.Key, tags, fields, time.Now()) - writer.WritePoint(p) + m.writePoint(writer, param.Key, fields, tags) } m.client.Close()