Skip to content
Closed
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
3 changes: 2 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
*.dll
*.so
*.dylib
**/.gocache

# Test binary, built with `go test -c`
*.test
Expand Down Expand Up @@ -39,4 +40,4 @@ opentelemetry-go

# Telegraf artifacts
config.toml
out.json
out.json**/.gocache
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
include ../../Makefile.Common
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
# Sketch Metrics Processor

This processor maintains Count-Min sketches over incoming metric datapoints (inspired by the Telegraf `countmin` aggregator). Each batch updates per-metric sketches keyed by tag dimensions, then emits a synthetic metric containing:

- Serialized Count-Min sketch payload (`countmin` attribute, bytes)
- Rows/columns/count metadata
- Approximate top-k heavy hitters (`topk` attribute, JSON)

Use `tag_keys` to explicitly choose which tags participate in the sketch. If omitted, all tags except those in `group_by` are used. Set `drop_original: true` to forward only the sketch metric.

See `config.go` for configuration details.
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package sketchcountminprocessor

import (
"context"
"testing"

"github.com/stretchr/testify/require"
"go.opentelemetry.io/collector/component/componenttest"
"go.opentelemetry.io/collector/consumer/consumertest"
"go.opentelemetry.io/collector/processor/processortest"
)

func TestComponentLifecycle(t *testing.T) {
factory := NewFactory()
cfg := factory.CreateDefaultConfig()

processor, err := factory.CreateMetrics(context.Background(), processortest.NewNopSettings(typeStr), cfg, consumertest.NewNop())
require.NoError(t, err)

host := componenttest.NewNopHost()
require.NoError(t, processor.Start(context.Background(), host))
require.NoError(t, processor.Shutdown(context.Background()))
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package sketchcountminprocessor // import "github.com/open-telemetry/opentelemetry-collector-contrib/processor/sketchmetricsprocessor"

import (
"fmt"

"go.opentelemetry.io/collector/component"
)

// Config defines configuration for the sketchmetrics processor.
//
// It is intentionally similar to the Telegraf countmin aggregator:
//
// [[aggregators.countmin]]
//
// measurement = "countmin"
// tag_keys = ["machineid", "tenant"]
// group_by = ["scrape_url"]
// rows = 3
// columns = 4096
// seed = 11400714819323198485
// top_k = 20
type Config struct {
// Name of the metric to emit (equivalent to Telegraf Measurement).
Measurement string `mapstructure:"measurement"`

// TagKeys explicitly specify which tag keys are used as the "by(...)" dimensions
// for top-k. If empty, all tags except GroupBy will be used.
TagKeys []string `mapstructure:"tag_keys"`

// GroupBy defines tags that partition the population into sub-groups.
// For each unique combination of group_by tags we maintain a separate set of sketches.
GroupBy []string `mapstructure:"group_by"`

// Count-Min sketch parameters.
Rows int `mapstructure:"rows"`
Columns int `mapstructure:"columns"`
Seed uint64 `mapstructure:"seed"`
TopK int `mapstructure:"top_k"`

// If true, original input metrics are dropped and only sketch metrics are forwarded.
DropOriginal bool `mapstructure:"drop_original"`
}

var _ component.Config = (*Config)(nil)

// Validate checks if the processor configuration is valid.
func (c *Config) Validate() error {
if c.TopK < 0 {
return fmt.Errorf("sketchmetrics: top_k must be >= 0")
}
if c.Rows < 0 || c.Columns < 0 {
return fmt.Errorf("sketchmetrics: rows/columns must be >= 0")
}
return nil
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package sketchcountminprocessor

import (
"testing"

"github.com/stretchr/testify/assert"
)

func TestConfigValidate(t *testing.T) {
cfg := Config{
Measurement: "countmin",
Rows: 3,
Columns: 4,
TopK: 10,
}
assert.NoError(t, cfg.Validate())

cfg.TopK = -1
assert.Error(t, cfg.Validate())
cfg.TopK = 1
cfg.Rows = -1
assert.Error(t, cfg.Validate())
cfg.Rows = 1
cfg.Columns = -1
assert.Error(t, cfg.Validate())
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package sketchcountminprocessor

import (
"bytes"
"encoding/binary"
"errors"
"fmt"
"hash"
"io"
"math"

"github.com/cespare/xxhash/v2"
)

const (
defaultRows = 3
defaultColumns = 4096
defaultSeed = 0x9e3779b185ebca87
)

var errInvalidConfig = errors.New("sketchmetrics: missing config")

type CountMinSketch struct {
rows int
cols int
total float64
table []float32
salts []uint64
hashes []hash.Hash64
topk *TopKHeap
}

func NewCountMinSketch(rows, cols int, seed uint64, topk int) (*CountMinSketch, error) {
if rows <= 0 || cols <= 0 {
return nil, fmt.Errorf("countmin: invalid dimensions rows=%d cols=%d", rows, cols)
}
cms := &CountMinSketch{
rows: rows,
cols: cols,
table: make([]float32, rows*cols),
salts: make([]uint64, rows),
hashes: make([]hash.Hash64, rows),
topk: NewTopKHeap(topk),
}
for i := 0; i < rows; i++ {
cms.salts[i] = mixSeed(seed, uint64(i))
cms.hashes[i] = xxhash.New()
}
return cms, nil
}

func (c *CountMinSketch) Insert(key string, weight float64) {
if weight == 0 {
return
}
estimate := math.MaxFloat64
for row := 0; row < c.rows; row++ {
h := c.hashes[row]
h.Reset()
var saltBuf [8]byte
binary.LittleEndian.PutUint64(saltBuf[:], c.salts[row])
h.Write(saltBuf[:])
io.WriteString(h, key)
idx := int(h.Sum64() % uint64(c.cols))
c.table[row*c.cols+idx] += float32(weight)
value := float64(c.table[row*c.cols+idx])
if value < estimate {
estimate = value
}
}
c.total += weight
if estimate == math.MaxFloat64 {
estimate = 0
}
if c.topk != nil {
c.topk.Update(key, estimate)
}
}

func (c *CountMinSketch) Rows() int {
return c.rows
}

func (c *CountMinSketch) Columns() int {
return c.cols
}

func (c *CountMinSketch) Total() float64 {
return c.total
}

func (c *CountMinSketch) TopKEntries() []TopKEntry {
if c.topk == nil {
return nil
}
return c.topk.Entries()
}

func (c *CountMinSketch) MarshalBinary() ([]byte, error) {
buf := &bytes.Buffer{}
if err := binary.Write(buf, binary.BigEndian, uint32(countMinMagic)); err != nil {
return nil, err
}
if err := binary.Write(buf, binary.BigEndian, uint16(countMinVersion)); err != nil {
return nil, err
}
if err := binary.Write(buf, binary.BigEndian, uint16(c.rows)); err != nil {
return nil, err
}
if err := binary.Write(buf, binary.BigEndian, uint32(c.cols)); err != nil {
return nil, err
}
if err := binary.Write(buf, binary.BigEndian, c.total); err != nil {
return nil, err
}
if err := binary.Write(buf, binary.BigEndian, uint16(len(c.salts))); err != nil {
return nil, err
}
for _, salt := range c.salts {
if err := binary.Write(buf, binary.BigEndian, salt); err != nil {
return nil, err
}
}
for _, v := range c.table {
if err := binary.Write(buf, binary.BigEndian, v); err != nil {
return nil, err
}
}
return buf.Bytes(), nil
}

const (
countMinMagic = 0x434d5331 // "CMS1"
countMinVersion = 1
)

func mixSeed(base uint64, row uint64) uint64 {
const prime uint64 = 0x100000001b3
value := base ^ (row * prime)
value ^= value >> 33
value *= 0xff51afd7ed558ccd
value ^= value >> 33
value *= 0xc4ceb9fe1a85ec53
value ^= value >> 33
return value
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

//go:generate mdatagen metadata.yaml

// Package sketchmetricsprocessor maintains Count-Min sketches over incoming metrics and
// emits serialized sketches and approximate top-k summaries.
package sketchcountminprocessor // import "github.com/open-telemetry/opentelemetry-collector-contrib/processor/sketchmetricsprocessor"
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package sketchcountminprocessor // import "github.com/open-telemetry/opentelemetry-collector-contrib/processor/sketchmetricsprocessor"

import (
"context"
"fmt"

"go.opentelemetry.io/collector/component"
"go.opentelemetry.io/collector/consumer"
"go.opentelemetry.io/collector/processor"
"go.opentelemetry.io/collector/processor/processorhelper"
)

var typeStr = component.MustNewType("sketchmetrics")

var capabilities = consumer.Capabilities{MutatesData: true}

// NewFactory returns a new factory for the sketchmetrics processor.
func NewFactory() processor.Factory {
return processor.NewFactory(
typeStr,
createDefaultConfig,
processor.WithMetrics(createMetricsProcessor, component.StabilityLevelAlpha),
)
}

func createDefaultConfig() component.Config {
return &Config{
Measurement: "countmin",
Rows: 3,
Columns: 4096,
Seed: 0x9e3779b185ebca87,
TopK: 20,
DropOriginal: false,
}
}

func createMetricsProcessor(
ctx context.Context,
set processor.Settings,
cfg component.Config,
next consumer.Metrics,
) (processor.Metrics, error) {
config, ok := cfg.(*Config)
if !ok {
return nil, fmt.Errorf("sketchmetrics: invalid config type %T", cfg)
}

p, err := newProcessor(config, set.Logger)
if err != nil {
return nil, err
}

return processorhelper.NewMetrics(
ctx,
set,
cfg,
next,
p.processMetrics,
processorhelper.WithCapabilities(capabilities),
)
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
package sketchcountminprocessor

import (
"testing"

"github.com/stretchr/testify/require"
"go.opentelemetry.io/collector/component/componenttest"
"go.opentelemetry.io/collector/consumer/consumertest"
"go.opentelemetry.io/collector/pipeline"
"go.opentelemetry.io/collector/processor/processortest"
)

func TestFactory_CreateDefaultConfig(t *testing.T) {
factory := NewFactory()
cfg := factory.CreateDefaultConfig()
require.NoError(t, componenttest.CheckConfigStruct(cfg))
}

func TestFactory_CreateMetrics(t *testing.T) {
factory := NewFactory()
cfg := factory.CreateDefaultConfig()

mp, err := factory.CreateMetrics(t.Context(), processortest.NewNopSettings(typeStr), cfg, consumertest.NewNop())
require.NoError(t, err)
require.NotNil(t, mp)
}

func TestFactory_CreateTracesNotSupported(t *testing.T) {
factory := NewFactory()
cfg := factory.CreateDefaultConfig()

tp, err := factory.CreateTraces(t.Context(), processortest.NewNopSettings(typeStr), cfg, consumertest.NewNop())
require.Nil(t, tp)
require.Equal(t, pipeline.ErrSignalNotSupported, err)
}
Loading