diff --git a/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-raw-metrics-flow.yaml b/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-raw-metrics-flow.yaml new file mode 100644 index 00000000..35510dc2 --- /dev/null +++ b/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-raw-metrics-flow.yaml @@ -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 diff --git a/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-sketch-payload-flow-batch.yaml b/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-sketch-payload-flow-batch.yaml new file mode 100644 index 00000000..154ef95a --- /dev/null +++ b/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-sketch-payload-flow-batch.yaml @@ -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 diff --git a/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-sketch-payload-flow-window.yaml b/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-sketch-payload-flow-window.yaml new file mode 100644 index 00000000..22857dd5 --- /dev/null +++ b/opentelemetry-collector-contrib-patch/cmd/kll/asap_query_config/config-otlp-sketch-payload-flow-window.yaml @@ -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 diff --git a/opentelemetry-collector-contrib-patch/cmd/kll/build-config.yaml b/opentelemetry-collector-contrib-patch/cmd/kll/build-config.yaml index d00cbaf0..ae9f5968 100644 --- a/opentelemetry-collector-contrib-patch/cmd/kll/build-config.yaml +++ b/opentelemetry-collector-contrib-patch/cmd/kll/build-config.yaml @@ -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 diff --git a/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.mod b/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.mod index 52d844e0..c8b492c6 100644 --- a/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.mod +++ b/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.mod @@ -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 ( @@ -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 ) @@ -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 diff --git a/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.sum b/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.sum index c42b4b6e..e6369c6b 100644 --- a/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.sum +++ b/opentelemetry-collector-contrib-patch/processor/kllprocessor/go.sum @@ -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= @@ -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= diff --git a/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor.go b/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor.go index a4fbc6ba..4baa3712 100644 --- a/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor.go +++ b/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor.go @@ -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 @@ -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 { @@ -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 { @@ -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 { diff --git a/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor_test.go b/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor_test.go index 0d39c10f..98cfcefb 100644 --- a/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor_test.go +++ b/opentelemetry-collector-contrib-patch/processor/kllprocessor/processor_test.go @@ -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) @@ -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)