Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions metric/distribution/distribution.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,10 @@ type Distribution interface {
Size() int

// weight is 1/samplingRate
AddEntryWithUnit(value float64, weight float64, unit string)
AddEntryWithUnit(value float64, weight float64, unit string) error

// weight is 1/samplingRate
AddEntry(value float64, weight float64)
AddEntry(value float64, weight float64) error

AddDistribution(distribution Distribution)

Expand Down
13 changes: 7 additions & 6 deletions metric/distribution/regular/regular_distribution.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
package regular

import (
"errors"
"log"
"math"

Expand Down Expand Up @@ -65,11 +66,10 @@ func (regularDist *RegularDistribution) Size() int {
}

// weight is 1/samplingRate
func (regularDist *RegularDistribution) AddEntryWithUnit(value float64, weight float64, unit string) {
func (regularDist *RegularDistribution) AddEntryWithUnit(value float64, weight float64, unit string) error {
if weight > 0 {
if value < 0 {
log.Printf("W! Value cannot be negative: %v", value)
return
return errors.New("negative value")
}
//sample count
regularDist.sampleCount += weight
Expand All @@ -91,16 +91,17 @@ func (regularDist *RegularDistribution) AddEntryWithUnit(value float64, weight f
if regularDist.unit == "" {
regularDist.unit = unit
} else if regularDist.unit != unit && unit != "" {
log.Printf("D! Multiple units are dected: %s, %s", regularDist.unit, unit)
log.Printf("D! Multiple units are detected: %s, %s", regularDist.unit, unit)
}
} else {
log.Printf("D! Weight should be larger than 0: %v", weight)
}
return nil
}

// weight is 1/samplingRate
func (regularDist *RegularDistribution) AddEntry(value float64, weight float64) {
regularDist.AddEntryWithUnit(value, weight, "")
func (regularDist *RegularDistribution) AddEntry(value float64, weight float64) error {
return regularDist.AddEntryWithUnit(value, weight, "")
}

func (regularDist *RegularDistribution) AddDistribution(distribution distribution.Distribution) {
Expand Down
6 changes: 3 additions & 3 deletions metric/distribution/regular/regular_distribution_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,9 @@ func TestSEH1Distribution(t *testing.T) {
//dist new and add entry
dist := NewRegularDistribution()

dist.AddEntry(20, 1)
dist.AddEntry(30, 1)
dist.AddEntryWithUnit(50, 1, "Count")
assert.NoError(t, dist.AddEntry(20, 1))
assert.NoError(t, dist.AddEntry(30, 1))
assert.NoError(t, dist.AddEntryWithUnit(50, 1, "Count"))

assert.Equal(t, 100.0, dist.Sum())
assert.Equal(t, 3.0, dist.SampleCount())
Expand Down
13 changes: 7 additions & 6 deletions metric/distribution/seh1/seh1_distribution.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
package seh1

import (
"errors"
"log"
"math"

Expand Down Expand Up @@ -75,11 +76,10 @@ func (seh1Distribution *SEH1Distribution) Size() int {
}

// weight is 1/samplingRate
func (seh1Distribution *SEH1Distribution) AddEntryWithUnit(value float64, weight float64, unit string) {
func (seh1Distribution *SEH1Distribution) AddEntryWithUnit(value float64, weight float64, unit string) error {
if weight > 0 {
if value < 0 {
log.Printf("W! Histogram value cannot be negative: %v", value)
return
return errors.New("negative value")
}
//sample count
seh1Distribution.sampleCount += weight
Expand All @@ -102,16 +102,17 @@ func (seh1Distribution *SEH1Distribution) AddEntryWithUnit(value float64, weight
if seh1Distribution.unit == "" {
seh1Distribution.unit = unit
} else if seh1Distribution.unit != unit && unit != "" {
log.Printf("D! Multiple units are dected: %s, %s", seh1Distribution.unit, unit)
log.Printf("D! Multiple units are detected: %s, %s", seh1Distribution.unit, unit)
}
} else {
log.Printf("D! Weight should be larger than 0: %v", weight)
}
return nil
}

// weight is 1/samplingRate
func (seh1Distribution *SEH1Distribution) AddEntry(value float64, weight float64) {
seh1Distribution.AddEntryWithUnit(value, weight, "")
func (seh1Distribution *SEH1Distribution) AddEntry(value float64, weight float64) error {
return seh1Distribution.AddEntryWithUnit(value, weight, "")
}

func (seh1Distribution *SEH1Distribution) AddDistribution(distribution distribution.Distribution) {
Expand Down
12 changes: 6 additions & 6 deletions metric/distribution/seh1/seh1_distribution_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,9 @@ func TestSEH1Distribution(t *testing.T) {
//dist new and add entry
dist := NewSEH1Distribution()

dist.AddEntry(20, 1)
dist.AddEntry(30, 1)
dist.AddEntryWithUnit(50, 1, "Count")
assert.NoError(t, dist.AddEntry(20, 1))
assert.NoError(t, dist.AddEntry(30, 1))
assert.NoError(t, dist.AddEntryWithUnit(50, 1, "Count"))

assert.Equal(t, 100.0, dist.Sum())
assert.Equal(t, 3.0, dist.SampleCount())
Expand All @@ -34,9 +34,9 @@ func TestSEH1Distribution(t *testing.T) {
//another dist new and add entry
anotherDist := NewSEH1Distribution()

anotherDist.AddEntry(21, 1)
anotherDist.AddEntry(22, 1)
anotherDist.AddEntry(23, 2)
assert.NoError(t, anotherDist.AddEntry(21, 1))
assert.NoError(t, anotherDist.AddEntry(22, 1))
assert.NoError(t, anotherDist.AddEntry(23, 2))

assert.Equal(t, 89.0, anotherDist.Sum())
assert.Equal(t, 4.0, anotherDist.SampleCount())
Expand Down
5 changes: 4 additions & 1 deletion plugins/inputs/statsd/statsd.go
Original file line number Diff line number Diff line change
Expand Up @@ -540,7 +540,10 @@ func (s *Statsd) aggregate(m metric) {
if m.samplerate > 0 {
weight = 1.0 / m.samplerate
}
field.(distribution.Distribution).AddEntry(m.floatvalue, weight)
err := field.(distribution.Distribution).AddEntry(m.floatvalue, weight)
if err != nil {
log.Printf("W! error: %s, metric: %s, value: %v", err, m.name, m.floatvalue)
}
cached.fields[m.field] = field
s.timings[m.hash] = cached
case "c":
Expand Down
50 changes: 25 additions & 25 deletions plugins/inputs/statsd/statsd_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -278,11 +278,11 @@ func TestParse_Timings(t *testing.T) {
s.Gather(acc)

dist := distribution.NewDistribution()
dist.AddEntry(1, 1)
dist.AddEntry(11, 1)
dist.AddEntry(1, 1)
dist.AddEntry(1, 1)
dist.AddEntry(1, 1)
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(11, 1))
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(1, 1))

metrics := acc.Metrics
assert.Equal(t, 1, len(metrics))
Expand Down Expand Up @@ -1039,18 +1039,18 @@ func TestParse_Timings_MultipleFieldsWithTemplate(t *testing.T) {
s.Gather(acc)

dist := distribution.NewDistribution()
dist.AddEntry(1, 1)
dist.AddEntry(11, 1)
dist.AddEntry(1, 1)
dist.AddEntry(1, 1)
dist.AddEntry(1, 1)
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(11, 1))
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(1, 1))

dist2 := distribution.NewDistribution()
dist2.AddEntry(2, 1)
dist2.AddEntry(22, 1)
dist2.AddEntry(2, 1)
dist2.AddEntry(2, 1)
dist2.AddEntry(2, 1)
assert.NoError(t, dist2.AddEntry(2, 1))
assert.NoError(t, dist2.AddEntry(22, 1))
assert.NoError(t, dist2.AddEntry(2, 1))
assert.NoError(t, dist2.AddEntry(2, 1))
assert.NoError(t, dist2.AddEntry(2, 1))

metrics := acc.Metrics
assert.Equal(t, 1, len(metrics))
Expand Down Expand Up @@ -1094,18 +1094,18 @@ func TestParse_Timings_MultipleFieldsWithoutTemplate(t *testing.T) {
s.Gather(acc)

dist := distribution.NewDistribution()
dist.AddEntry(1, 1)
dist.AddEntry(11, 1)
dist.AddEntry(1, 1)
dist.AddEntry(1, 1)
dist.AddEntry(1, 1)
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(11, 1))
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(1, 1))
assert.NoError(t, dist.AddEntry(1, 1))

dist2 := distribution.NewDistribution()
dist2.AddEntry(2, 1)
dist2.AddEntry(22, 1)
dist2.AddEntry(2, 1)
dist2.AddEntry(2, 1)
dist2.AddEntry(2, 1)
assert.NoError(t, dist2.AddEntry(2, 1))
assert.NoError(t, dist2.AddEntry(22, 1))
assert.NoError(t, dist2.AddEntry(2, 1))
assert.NoError(t, dist2.AddEntry(2, 1))
assert.NoError(t, dist2.AddEntry(2, 1))

metrics := acc.Metrics
assert.Equal(t, 2, len(metrics))
Expand Down
5 changes: 4 additions & 1 deletion plugins/outputs/cloudwatch/aggregator.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,10 @@ func (durationAgg *durationAggregator) aggregating() {
if dist != nil {
existingDist.AddDistribution(dist)
} else {
existingDist.AddEntry(value, 1)
err = existingDist.AddEntry(value, 1)
if err != nil {
log.Printf("W! error: %s, metric %s, value %v", err, m.Name(), value)
}
}
}
case <-durationAgg.ticker.C:
Expand Down
4 changes: 3 additions & 1 deletion plugins/outputs/cloudwatch/cloudwatch_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,9 @@ func TestBuildMetricDatums(t *testing.T) {
}

invalidDistribution := distribution.NewDistribution()
invalidDistribution.AddEntry(-1, 1)
err := invalidDistribution.AddEntry(-1, 1)
expectedErrMsg := "negative value"
assert.EqualError(err, expectedErrMsg)
invalidMetrics := []telegraf.Metric{
testutil.TestMetric("Foo"),
testutil.TestMetric(invalidDistribution),
Expand Down
6 changes: 3 additions & 3 deletions plugins/outputs/cloudwatch/util_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,9 @@ func TestResize(t *testing.T) {
assert.Equal(t, float64(1), sampleCount)
assert.Equal(t, float64(1), sum)

dist.AddEntry(2, 1)
dist.AddEntry(3, 1)
dist.AddEntry(4, 1)
assert.NoError(t, dist.AddEntry(2, 1))
assert.NoError(t, dist.AddEntry(3, 1))
assert.NoError(t, dist.AddEntry(4, 1))

distList = resize(dist, maxListSize)
assert.Equal(t, 2, len(distList))
Expand Down