InfluxDB: write slice of structs (#5873)
This commit is contained in:
parent
29bdd1a7ac
commit
0f07cd612c
1 changed files with 50 additions and 43 deletions
|
|
@ -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()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue