Limit influxdb writes to minute resolution and don't write duplicate socCharge values
This commit is contained in:
parent
871768e3b5
commit
fb04aa98da
3 changed files with 119 additions and 1 deletions
11
cmd/root.go
11
cmd/root.go
|
|
@ -135,6 +135,8 @@ func checkVersion() {
|
|||
|
||||
var teeIsChained bool // controles piping of first channel in teed chain
|
||||
|
||||
// tee splits a tee channel from an input channel and starts a goroutine
|
||||
// for duplicateing channel values to replacement input and tee channel
|
||||
func tee(in chan core.Param) (chan core.Param, <-chan core.Param) {
|
||||
gen := make(chan core.Param)
|
||||
tee := make(chan core.Param)
|
||||
|
|
@ -205,6 +207,14 @@ func run(cmd *cobra.Command, args []string) {
|
|||
var teeChan <-chan core.Param
|
||||
valueChan, teeChan = tee(valueChan)
|
||||
|
||||
// eliminate duplicate values
|
||||
dedupe := server.NewDeduplicator(30*time.Minute, "socCharge")
|
||||
teeChan = dedupe.Pipe(teeChan)
|
||||
|
||||
// reduce number of values written to influx
|
||||
limiter := server.NewLimiter(1 * time.Minute)
|
||||
teeChan = limiter.Pipe(teeChan)
|
||||
|
||||
go influx.Run(teeChan)
|
||||
}
|
||||
|
||||
|
|
@ -214,6 +224,7 @@ func run(cmd *cobra.Command, args []string) {
|
|||
|
||||
var teeChan <-chan core.Param
|
||||
valueChan, teeChan = tee(valueChan)
|
||||
|
||||
go socketHub.Run(teeChan, triggerChan)
|
||||
|
||||
// start all loadpoints
|
||||
|
|
|
|||
|
|
@ -1,8 +1,16 @@
|
|||
uri: 0.0.0.0:7070 # uri for ui
|
||||
interval: 10s # control cycle interval
|
||||
|
||||
# mqtt message broker
|
||||
mqtt:
|
||||
broker: localhost:1883 # mqtt broker address
|
||||
broker: localhost:1883
|
||||
# user:
|
||||
# password:
|
||||
|
||||
# influx database (v1)
|
||||
influx:
|
||||
url: http://localhost:8086
|
||||
database: evcc
|
||||
# user:
|
||||
# password:
|
||||
|
||||
|
|
|
|||
99
server/limiter.go
Normal file
99
server/limiter.go
Normal file
|
|
@ -0,0 +1,99 @@
|
|||
package server
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"github.com/andig/evcc/core"
|
||||
)
|
||||
|
||||
// Piper is the interface that data flow plugins must implement
|
||||
type Piper interface {
|
||||
Pipe(in <-chan core.Param) <-chan core.Param
|
||||
}
|
||||
|
||||
type cacheItem struct {
|
||||
updated time.Time
|
||||
val interface{}
|
||||
}
|
||||
|
||||
// Deduplicator allows filtering of channel data by given criteria
|
||||
type Deduplicator struct {
|
||||
interval time.Duration
|
||||
filter map[string]interface{}
|
||||
cache map[string]cacheItem
|
||||
}
|
||||
|
||||
// NewDeduplicator creates Deduplicator
|
||||
func NewDeduplicator(interval time.Duration, filter ...string) Piper {
|
||||
l := &Deduplicator{
|
||||
interval: interval,
|
||||
filter: make(map[string]interface{}),
|
||||
cache: make(map[string]cacheItem),
|
||||
}
|
||||
|
||||
for _, f := range filter {
|
||||
l.filter[f] = struct{}{}
|
||||
}
|
||||
|
||||
return l
|
||||
}
|
||||
|
||||
func (l *Deduplicator) pipe(in <-chan core.Param, out chan<- core.Param) {
|
||||
for p := range in {
|
||||
// use loadpoint + param.Key as lookup key to value cache
|
||||
key := p.LoadPoint + "." + p.Key
|
||||
item, cached := l.cache[key]
|
||||
_, filtered := l.filter[p.Key]
|
||||
|
||||
// forward if not cached
|
||||
if !cached || !filtered || filtered &&
|
||||
(time.Since(item.updated) >= l.interval || p.Val != item.val) {
|
||||
l.cache[key] = cacheItem{updated: time.Now(), val: p.Val}
|
||||
out <- p
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pipe creates a new filtered output channel for given input channel
|
||||
func (l *Deduplicator) Pipe(in <-chan core.Param) <-chan core.Param {
|
||||
out := make(chan core.Param)
|
||||
go l.pipe(in, out)
|
||||
return out
|
||||
}
|
||||
|
||||
// Limiter allows filtering of channel data by given criteria
|
||||
type Limiter struct {
|
||||
interval time.Duration
|
||||
cache map[string]cacheItem
|
||||
}
|
||||
|
||||
// NewLimiter creates limiter
|
||||
func NewLimiter(interval time.Duration) Piper {
|
||||
l := &Limiter{
|
||||
interval: interval,
|
||||
cache: make(map[string]cacheItem),
|
||||
}
|
||||
|
||||
return l
|
||||
}
|
||||
|
||||
func (l *Limiter) pipe(in <-chan core.Param, out chan<- core.Param) {
|
||||
for p := range in {
|
||||
// use loadpoint + param.Key as lookup key to value cache
|
||||
key := p.LoadPoint + "." + p.Key
|
||||
item, cached := l.cache[key]
|
||||
|
||||
// forward if not cached or expired
|
||||
if !cached || time.Since(item.updated) >= l.interval {
|
||||
l.cache[key] = cacheItem{updated: time.Now(), val: p.Val}
|
||||
out <- p
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pipe creates a new filtered output channel for given input channel
|
||||
func (l *Limiter) Pipe(in <-chan core.Param) <-chan core.Param {
|
||||
out := make(chan core.Param)
|
||||
go l.pipe(in, out)
|
||||
return out
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue