From cb3141ff44ba144e4302aa63d28a581252fca604 Mon Sep 17 00:00:00 2001 From: Zeying Zhu Date: Wed, 15 Apr 2026 16:38:49 -0400 Subject: [PATCH] feat(countminsketchprocessor): emit typed CountMinSketchDataPoint MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Refactors the TransmitSketch path from Gauge-with-byte-attribute emission to typed `CountMinSketchDataPoint` messages so ASAPQuery- backend's modified-OTLP sketch router (ASAPQuery-backend PRs #5-#9) actually sees them as sketch variants instead of anonymous Gauges. ## Why Before this PR the processor emitted: metric.SetEmptyGauge() dp := gauge.DataPoints().AppendEmpty() dp.Attributes().PutEmptyBytes("sketch_payload").FromRaw(payload) dp.Attributes().PutStr("encoding", "proto_full") ASAPQuery-backend's modified-OTLP decoder matches on `Metric.data = CountMinSketch{data_points: [...]}` (oneof tag 16, typed `CountMinSketchDataPoint` messages) and never looks inside anonymous Gauge attribute maps for sketch bytes. Even though the backend has a full CountMin decoder (PR #6), the whole end-to-end flow was broken for this processor because nothing produced the typed data points the decoder wanted. After this PR: metric.SetEmptyCountMinSketch() cmsMetric.SetAggregationTemporality(pmetric.AggregationTemporalityDelta) dp := cmsMetric.DataPoints().AppendEmpty() outputAttrs.CopyTo(dp.Attributes()) dp.SetSampleCount(...) dp.SetRows(...) dp.SetCols(...) dp.SetSketch(payload) dp.SetEncoding(pmetric.CountMinSketchEncodingProto | Delta) The backend's router now sees `Metric.data.CountMinSketch{...}` and routes it straight into `CountMinSketchAccumulator::from_sketchlib_proto_bytes`. This unblocks the ASAPQuery PR #11 / #12 feedback loop for CountMin queries — a capability miss now leads to a plan push which leads to the backend actually finding a precomputed match next time, because the sketch bytes are now flowing through the typed hot path. ## What changed ### `processor.go` Split the emission path on `p.cfg.TransmitSketch`: * **TransmitSketch = true** (the production path): emit a typed `CountMinSketchDataPoint` with `SetSketch` / `SetEncoding` / `SetSampleCount` / `SetRows` / `SetCols`. The internal encoding string ("proto_delta" / "proto_full") is mapped onto the proto enum (`CountMinSketchEncodingDelta` / `...Proto`). * **TransmitSketch = false** (monitoring-only path): keep the legacy Gauge emission so existing dashboards that read the processor's output as a scalar `sample_count` series keep working. Nothing in that mode carries sketch bytes anyway. `aggregationTemporality` is set to Delta because the processor always produces one data point per window, and every window represents the delta within that window (not a cumulative sketch). The backend's accumulator merges across windows itself. ### `processor_test.go` Introduced a tiny `cmsTestDataPoint` adapter so existing test assertions that read fields via `.Attributes().Get("sketch_payload")` / `"encoding"` / `"sample_count"` / `"rows"` / `"cols"` keep compiling without per-site rewrites. The helper `getAllDataPoints` now dispatches on `pmetric.MetricTypeCountMinSketch` for the typed path and synthesizes the legacy attribute keys from the typed fields (`Sketch()` / `Encoding()` / etc.). The Gauge path (`TransmitSketch = false`) still works via the same helper — the adapter carries `doubleValue` for the one test (`TestBatchModeQueryMetricsWhenTransmitSketchDisabled`) that reads `dps[0].DoubleValue()`. `encodingToLegacyString` maps the proto enum back to the "proto_full" / "proto_delta" strings that existing assertions compare against. Drive-by fix: two assertions in processor_test.go previously looked for `"cms.sketch_payload"` (with a `cms.` prefix) which had never matched what the processor actually wrote (`"sketch_payload"`, no prefix). Those assertions were silently always failing on any CI that exercised them. Fixed to `"sketch_payload"` to match the (legacy-compatible, now synthesized) attribute key. ## What's NOT in this PR * **MSGPACK encoding option**: the typed emission path is now ready to accept a config flag (next PR will add `encoding: msgpack` and call sketchlib-go's `SerializeMsgpack` from PR #51, setting `CountMinSketchEncodingMsgpack` from PR #157). One-line branch once this refactor lands. * **The other 3 processors** (`countsketchprocessor`, `hllprocessor`, `kllprocessor`) need the same refactor but each has its own attribute quirks and test suite. Will ship as 3 separate PRs to keep review size bounded. ## Validation Local `go build` was attempted but the worktree's processor module is missing a `replace` directive for `go.opentelemetry.io/collector/processor`, so `processor/selfmonitor` fails to resolve — a pre-existing build setup issue unrelated to this PR. `gofmt -l` is clean on both modified files. All my changes use already-existing pmetric API methods (`SetEmptyCountMinSketch`, `SetSketch`, `SetEncoding`, `SetSampleCount`, `SetRows`, `SetCols`) and existing enum constants (`CountMinSketchEncodingProto` / `Delta`), so compilation should be straightforward in CI's properly-set-up module graph. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../countminsketchprocessor/processor.go | 49 +++++++++--- .../countminsketchprocessor/processor_test.go | 77 +++++++++++++++++-- 2 files changed, 109 insertions(+), 17 deletions(-) diff --git a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go index dc6c059f..9d2e57e7 100644 --- a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go +++ b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go @@ -520,18 +520,47 @@ func (p *windowedCountMinSketchProcessor) buildWindowMetricsAndReset() pmetric.M m.SetName(p.cfg.MetricName) m.SetUnit("1") - gauge := m.SetEmptyGauge() - dp := gauge.DataPoints().AppendEmpty() - dp.SetTimestamp(now) - - outputAttrs.CopyTo(dp.Attributes()) - dp.Attributes().PutInt("rows", int64(rows)) - dp.Attributes().PutInt("cols", int64(cols)) - dp.Attributes().PutInt("sample_count", int64(sampleCount)) if p.cfg.TransmitSketch { - dp.Attributes().PutEmptyBytes("sketch_payload").FromRaw(payload) - dp.Attributes().PutStr("encoding", encoding) + // Typed CountMinSketchDataPoint emission — what + // ASAPQuery-backend's modified-OTLP sketch router + // consumes as `Metric.data = CountMinSketch{...}`. + // Before this change the processor emitted a Gauge + // with the sketch payload stuffed into a + // `sketch_payload` byte attribute, which the backend + // router never recognized as a sketch variant. + cmsMetric := m.SetEmptyCountMinSketch() + cmsMetric.SetAggregationTemporality(pmetric.AggregationTemporalityDelta) + dp := cmsMetric.DataPoints().AppendEmpty() + dp.SetTimestamp(now) + outputAttrs.CopyTo(dp.Attributes()) + dp.SetSampleCount(uint64(sampleCount)) + dp.SetRows(int32(rows)) + dp.SetCols(int32(cols)) + dp.SetSketch(payload) + // Map the internal encoding string onto the proto + // enum the backend expects. The encoding string only + // branches on delta vs full when delta transmission + // is on; otherwise it's always proto_full. + switch encoding { + case "proto_delta": + dp.SetEncoding(pmetric.CountMinSketchEncodingDelta) + default: + // "proto_full" and any unexpected fallback. + dp.SetEncoding(pmetric.CountMinSketchEncodingProto) + } } else { + // Non-transmit mode: caller only wants the + // per-window sample count for monitoring, not the + // sketch bytes. Keep the legacy Gauge emission so + // existing dashboards that read `countmin` as a + // scalar series continue to work. + gauge := m.SetEmptyGauge() + dp := gauge.DataPoints().AppendEmpty() + dp.SetTimestamp(now) + outputAttrs.CopyTo(dp.Attributes()) + dp.Attributes().PutInt("rows", int64(rows)) + dp.Attributes().PutInt("cols", int64(cols)) + dp.Attributes().PutInt("sample_count", int64(sampleCount)) dp.SetDoubleValue(float64(sampleCount)) } } diff --git a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go index cf31e918..5a2f57d3 100644 --- a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go +++ b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go @@ -15,6 +15,7 @@ import ( "github.com/stretchr/testify/require" "go.opentelemetry.io/collector/component/componenttest" "go.opentelemetry.io/collector/consumer" + "go.opentelemetry.io/collector/pdata/pcommon" "go.opentelemetry.io/collector/pdata/pmetric" "go.uber.org/zap" ) @@ -180,7 +181,7 @@ func TestProcessor_TumblingWindow_Correctness(t *testing.T) { // ========================================== // 3. Verify Binary Payload (Gob Decode) // ========================================== - payloadVal, ok := dps2[0].Attributes().Get("cms.sketch_payload") + payloadVal, ok := dps2[0].Attributes().Get("sketch_payload") require.True(t, ok, "Sketch payload must exist in attributes") rawBytes := payloadVal.Bytes().AsRaw() @@ -210,8 +211,45 @@ func generateMetrics(serviceName string, count int) pmetric.Metrics { return md } -func getAllDataPoints(md pmetric.Metrics) []pmetric.NumberDataPoint { - var dps []pmetric.NumberDataPoint +// cmsTestDataPoint is a test-only adapter that flattens either a +// `CountMinSketchDataPoint` (the typed emission path) or a `NumberDataPoint` +// (the legacy Gauge path kept for non-TransmitSketch mode) into a single +// `pcommon.Map` so existing test assertions that read `sketch_payload`, +// `encoding`, `sample_count`, `rows`, `cols` via `Attributes().Get(...)` +// continue to compile without site-by-site rewrites. +// +// `doubleValue` is only populated for the non-TransmitSketch Gauge path — +// `TestBatchModeQueryMetricsWhenTransmitSketchDisabled` is the single +// test that reads it. For the typed-DP path it stays 0 and is unused. +type cmsTestDataPoint struct { + attributes pcommon.Map + doubleValue float64 +} + +func (d cmsTestDataPoint) Attributes() pcommon.Map { return d.attributes } +func (d cmsTestDataPoint) DoubleValue() float64 { return d.doubleValue } + +// encodingToLegacyString mirrors the strings the processor used to +// write into the `encoding` attribute, so tests that compare against +// "proto_full" / "proto_delta" continue to work. +func encodingToLegacyString(enc pmetric.CountMinSketchEncoding) string { + switch enc { + case pmetric.CountMinSketchEncodingProto: + return "proto_full" + case pmetric.CountMinSketchEncodingDelta: + return "proto_delta" + } + return "unknown" +} + +// getAllDataPoints walks the output metrics and returns a flat slice of +// test adapters, one per sketch data point (typed or legacy Gauge). When +// the input metric is a typed `CountMinSketch`, its fields (sketch bytes, +// encoding enum, sample_count, rows, cols) are synthesized into the +// adapter's attribute map under their legacy string keys so the rest of +// the test file can keep reading them with `Attributes().Get(...)`. +func getAllDataPoints(md pmetric.Metrics) []cmsTestDataPoint { + var dps []cmsTestDataPoint rms := md.ResourceMetrics() for i := 0; i < rms.Len(); i++ { sms := rms.At(i).ScopeMetrics() @@ -219,9 +257,34 @@ func getAllDataPoints(md pmetric.Metrics) []pmetric.NumberDataPoint { ms := sms.At(j).Metrics() for k := 0; k < ms.Len(); k++ { m := ms.At(k) - pts := m.Gauge().DataPoints() - for l := 0; l < pts.Len(); l++ { - dps = append(dps, pts.At(l)) + switch m.Type() { + case pmetric.MetricTypeCountMinSketch: + pts := m.CountMinSketch().DataPoints() + for l := 0; l < pts.Len(); l++ { + dp := pts.At(l) + attrs := pcommon.NewMap() + dp.Attributes().CopyTo(attrs) + // Inject the typed-DP fields under their + // legacy attribute names so existing tests + // keep working. + attrs.PutEmptyBytes("sketch_payload").FromRaw(dp.Sketch()) + attrs.PutStr("encoding", encodingToLegacyString(dp.Encoding())) + attrs.PutInt("sample_count", int64(dp.SampleCount())) + attrs.PutInt("rows", int64(dp.Rows())) + attrs.PutInt("cols", int64(dp.Cols())) + dps = append(dps, cmsTestDataPoint{attributes: attrs}) + } + case pmetric.MetricTypeGauge: + pts := m.Gauge().DataPoints() + for l := 0; l < pts.Len(); l++ { + dp := pts.At(l) + attrs := pcommon.NewMap() + dp.Attributes().CopyTo(attrs) + dps = append(dps, cmsTestDataPoint{ + attributes: attrs, + doubleValue: dp.DoubleValue(), + }) + } } } } @@ -301,7 +364,7 @@ func TestBatchModeQueryMetricsWhenTransmitSketchDisabled(t *testing.T) { dps := getAllDataPoints(out) require.Len(t, dps, 1) assert.Equal(t, 3.0, dps[0].DoubleValue()) - _, hasPayload := dps[0].Attributes().Get("cms.sketch_payload") + _, hasPayload := dps[0].Attributes().Get("sketch_payload") assert.False(t, hasPayload) }