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
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
# OTLP Raw Metrics Flow - Local testing
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:53217
http:
endpoint: 0.0.0.0:53218

processors:
batch:
send_batch_size: 5000

exporters:
otlphttp:
endpoint: http://localhost:4318
encoding: proto
compression: gzip
service:
pipelines:
metrics:
receivers: [otlp]
processors: [batch]
exporters: [otlphttp]
telemetry:
metrics:
level: basic
logs:
level: info
resource:
service.name: kllcol-otlp-raw-metrics-flow-local
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
# OTLP Sketch Payload Flow - Batch mode - Local testing
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:53217
http:
endpoint: 0.0.0.0:53218

processors:
KLL:
mode: batch
k: 256
transmit_sketch: true
quantiles: [0.5, 0.9, 0.99]
drop_original: true

exporters:
otlphttp:
endpoint: http://localhost:4318
encoding: proto
compression: gzip
service:
pipelines:
metrics:
receivers: [otlp]
processors: [KLL]
exporters: [otlphttp]
telemetry:
metrics:
level: basic
logs:
level: debug
resource:
service.name: kllcol-otlp-sketch-payload-flow-local
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
# OTLP Sketch Payload Flow - Window mode - Local testing
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:53217
http:
endpoint: 0.0.0.0:53218

processors:
KLL:
mode: window
window_duration: 10s
k: 256
transmit_sketch: true
quantiles: [0.5, 0.9, 0.99]
drop_original: true

exporters:
otlphttp:
endpoint: http://localhost:4318
encoding: proto
compression: gzip
service:
pipelines:
metrics:
receivers: [otlp]
processors: [KLL]
exporters: [otlphttp]
telemetry:
metrics:
level: basic
logs:
level: debug
resource:
service.name: kllcol-otlp-sketch-payload-flow-window-local
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,19 @@ dist:
output_path: .

exporters:
- gomod: go.opentelemetry.io/collector/exporter/nopexporter v0.141.0
- gomod: go.opentelemetry.io/collector/exporter/nopexporter v0.142.0
- gomod: go.opentelemetry.io/collector/exporter/otlphttpexporter v0.142.0
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/exporter/fileexporter v0.142.0
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/exporter/prometheusexporter v0.142.0

processors:
- gomod: go.opentelemetry.io/collector/processor/batchprocessor v0.141.0
- gomod: go.opentelemetry.io/collector/processor/batchprocessor v0.142.0
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/processor/kllprocessor v0.0.0

receivers:
- gomod: go.opentelemetry.io/collector/receiver/otlpreceiver v0.141.0
- gomod: go.opentelemetry.io/collector/receiver/otlpreceiver v0.142.0

replaces:
- github.com/open-telemetry/opentelemetry-collector-contrib/processor/kllprocessor => ./processor/kllprocessor
- go.opentelemetry.io/collector/pdata => ../opentelemetry-collector/pdata
# - github.com/ProjectASAP/sketchlib-go => ../../sketchlib-go
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,13 @@ module github.com/open-telemetry/opentelemetry-collector-contrib/processor/kllpr
go 1.25.4

require (
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321024028-d20a9f9151b5
github.com/ProjectASAP/sketchlib-go v0.0.0-20260310013347-2c5db0c75da8
github.com/stretchr/testify v1.11.1
go.opentelemetry.io/collector/component v1.48.0
go.opentelemetry.io/collector/component/componenttest v0.142.0
go.opentelemetry.io/collector/consumer v1.48.0
go.opentelemetry.io/collector/consumer/consumertest v0.142.0
google.golang.org/protobuf v1.36.11
)

require (
Expand All @@ -32,7 +33,6 @@ require (
go.yaml.in/yaml/v2 v2.4.3 // indirect
golang.org/x/sys v0.39.0 // indirect
golang.org/x/text v0.30.0 // indirect
google.golang.org/protobuf v1.36.11 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)

Expand All @@ -54,4 +54,4 @@ require (

replace go.opentelemetry.io/collector/pdata => ../../../opentelemetry-collector/pdata

replace go.opentelemetry.io/collector/processor => ../../../opentelemetry-collector/processor
// replace github.com/ProjectASAP/sketchlib-go => ../../../../sketchlib-go
Original file line number Diff line number Diff line change
@@ -1,13 +1,3 @@
github.com/ProjectASAP/sketchlib-go v0.0.0-20260316040945-49890c71f035 h1:aLrGKic2Q7KJJmfTzqB3fSvjjW2KfRvY5BDcRvYI2zE=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260316040945-49890c71f035/go.mod h1:Nzyu+hI1KUBmgeyhSlmUnTp5LvI/EwD8VuoZyNQoEZs=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260320220729-3ba826ceb054 h1:HYCqQ0IEu/hDk5UOi4AHF5yuNb/DX3jD8x2PtO4bApQ=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260320220729-3ba826ceb054/go.mod h1:VmV0RYT6+rXpb1gVxTG3Z17sttm0nfj7RuesrL1oDTI=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321021603-cddd774cb224 h1:mYN1NJLRp9elj4RoLqlHMcoZeQxI7jPNqnxbxxS7wBw=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321021603-cddd774cb224/go.mod h1:VmV0RYT6+rXpb1gVxTG3Z17sttm0nfj7RuesrL1oDTI=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321023259-ecebc36fb5aa h1:nRWz+s995NYgaRI6a7YHu2lnfGhyzn3Sn8uDFbuNEO0=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321023259-ecebc36fb5aa/go.mod h1:VmV0RYT6+rXpb1gVxTG3Z17sttm0nfj7RuesrL1oDTI=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321024028-d20a9f9151b5 h1:e1If3BxB2QRDt/57ho0PoxxIwouEvuaojy8m9u21IC0=
github.com/ProjectASAP/sketchlib-go v0.0.0-20260321024028-d20a9f9151b5/go.mod h1:VmV0RYT6+rXpb1gVxTG3Z17sttm0nfj7RuesrL1oDTI=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
Expand Down Expand Up @@ -117,8 +107,6 @@ golang.org/x/sys v0.39.0 h1:CvCKL8MeisomCi6qNZ+wbb0DN9E5AATixKsvNtMoMFk=
golang.org/x/sys v0.39.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
golang.org/x/text v0.30.0 h1:yznKA/E9zq54KzlzBEAWn1NXSQ8DIp/NYMy88xJjl4k=
golang.org/x/text v0.30.0/go.mod h1:yDdHFIX9t+tORqspjENWgzaCVXgk0yYnYuSZ8UzzBVM=
google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE=
google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"go.opentelemetry.io/collector/pdata/pmetric"
"go.opentelemetry.io/collector/processor/selfmonitor"
"go.uber.org/zap"
"google.golang.org/protobuf/proto"
)

// builderPool recycles strings.Builder instances to avoid per-call heap
Expand Down Expand Up @@ -197,7 +198,7 @@ func (p *kllProcessor) processBatch(md pmetric.Metrics) error {
}
bs := getOrCreate(metric.Name(), metric.Unit(), dp.Attributes())
if bs.sketch != nil && len(dp.Sketch()) > 0 {
incoming, err := kll.DeserializeKLLSketchFromProtoBytes(dp.Sketch())
incoming, err := kll.DeserializeKLLSketchFromBytes(dp.Sketch())
if err == nil {
_ = bs.sketch.Merge(incoming)
} else if p.logger != nil {
Expand Down Expand Up @@ -440,7 +441,7 @@ func (p *kllProcessor) accumulateKLLSketchMetric(sw *scopeWindow, metric pmetric
mw.series[attrKey] = series
}
if series.sketch != nil && len(dp.Sketch()) > 0 {
incoming, err := kll.DeserializeKLLSketchFromProtoBytes(dp.Sketch())
incoming, err := kll.DeserializeKLLSketchFromBytes(dp.Sketch())
if err == nil {
_ = series.sketch.Merge(incoming)
} else if p.logger != nil {
Expand Down Expand Up @@ -629,7 +630,11 @@ func serializeKLLSketch(sketch *kll.KLLSketch) ([]byte, error) {
if sketch == nil {
return nil, nil
}
return sketch.SerializeProtoBytes()
env, err := sketch.SerializePortable()
if err != nil {
return nil, err
}
return proto.Marshal(env)
}

func (p *kllProcessor) sketchMetricName(base string) string {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ func TestBatchModeTransmitSketch(t *testing.T) {
cfg.Mode = ModeBatch
cfg.TransmitSketch = true
cfg.Quantiles = nil
cfg.DropOriginal = true
require.NoError(t, cfg.Validate())

sink := new(consumertest.MetricsSink)
Expand Down Expand Up @@ -290,10 +291,12 @@ func TestEmptyInput(t *testing.T) {
}

// TestEmptyResourceMetrics verifies ResourceMetrics with zero ScopeMetrics is handled in batch mode.
// Uses DropOriginal=false to test passthrough when there is no data to aggregate.
func TestEmptyResourceMetrics(t *testing.T) {
cfg := createDefaultConfig().(*Config)
cfg.Mode = ModeBatch
cfg.Quantiles = []float64{0.5}
cfg.DropOriginal = false
require.NoError(t, cfg.Validate())

sink := new(consumertest.MetricsSink)
Expand Down