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
6 changes: 3 additions & 3 deletions controller/src/config/stage_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -439,7 +439,7 @@ pub fn emit_gateway_yaml(cfg: &GatewayStageConfig, opamp_endpoint: &str) -> Resu
// * SketchKind::DDSketch → `ddsketchmerge`
// * SketchKind::Kll → `kllmerge`
// * SketchKind::Hll → `hllmerge`
// * SketchKind::Cms → `countminmerge`
// * SketchKind::Cms → `countminsketchmerge`
// * SketchKind::CountSketch → `countsketchmerge`
//
// We honour `GatewayMergeProcessor::processor_name` if non-empty
Expand Down Expand Up @@ -1320,15 +1320,15 @@ fn build_edge_processor_block(
/// Today the typed emitter populates every entry's `processor_name`
/// with the placeholder `"sketchmergeprocessor"`; the patched contrib
/// build instead has per-family merge processors:
/// `kllmerge`, `ddsketchmerge`, `hllmerge`, `countminmerge`,
/// `kllmerge`, `ddsketchmerge`, `hllmerge`, `countminsketchmerge`,
/// `countsketchmerge`. We map the kind to the family-specific name
/// here so the emitted YAML round-trips through the patched build.
fn gateway_merge_processor_name(mp: &GatewayMergeProcessor) -> String {
match mp.sketch_kind {
SketchKind::Kll => "kllmerge".to_string(),
SketchKind::DDSketch => "ddsketchmerge".to_string(),
SketchKind::Hll => "hllmerge".to_string(),
SketchKind::Cms => "countminmerge".to_string(),
SketchKind::Cms => "countminsketchmerge".to_string(),
SketchKind::CountSketch => "countsketchmerge".to_string(),
}
}
Expand Down
32 changes: 20 additions & 12 deletions deploy/scripts/run_mvp_demo.sh
Original file line number Diff line number Diff line change
Expand Up @@ -427,10 +427,16 @@ capture_emitted_configs() {
fi

# Per-metric typed config (one per workload entry — Phase B
# emitter output). The mvp-workload.yaml has four entries.
# emitter output). mvp-workload.yaml has 6 distinct metrics
# across 8 entries (http_requests_total appears in 3 query
# shapes: raw passthrough, gateway sum, cold archive probe).
for metric in \
http_requests_total_latency_ms \
http_requests_total ; do
http_requests_total \
request_size_bytes \
unique_users_per_min \
top_endpoint_qps \
endpoint_request_freq ; do
local out="${cdir}/per-metric.${metric}.json"
if curl -sf "${ctrl}/api/v1/config/${metric}" -o "${out}" \
2> "${out}.err"; then
Expand Down Expand Up @@ -498,18 +504,20 @@ measure_phase() {

# Build the replay query suite from mvp-workload.yaml.
# We keep the JSON adjacent to the run dir for reproducibility.
# Six query classes — one per sketch family registered in
# mvp-workload.yaml (issue #46 5-sketch coverage):
# Seven query classes — one per sketch family + the label-agg
# family appears in two shapes (issue #46 5-sketch + 3 canonical
# query classes):
#
# sum_rate ↔ raw passthrough (http_requests_total)
# quantile ↔ DDSketch (http_requests_total_latency_ms)
# kll-quantile ↔ KLL (request_size_bytes)
# count_unique ↔ HLL (unique_users_per_min)
# topk ↔ CountSketch (top_endpoint_qps)
# frequency ↔ CountMinSketch (endpoint_request_freq)
# quantile ↔ DDSketch (http_requests_total_latency_ms)
# sum-by-zone ↔ raw passthrough (http_requests_total) [criterion ① label-at-instant]
# combined-rate ↔ raw passthrough (http_requests_total) [criterion ① combined window+label]
# kll-quantile ↔ KLL (request_size_bytes)
# count_unique ↔ HLL (unique_users_per_min)
# topk ↔ CountSketch (top_endpoint_qps)
# frequency ↔ CountMinSketch (endpoint_request_freq)
#
# Replay rotates round-robin at QPS=8 → ~1.3 QPS per class →
# ≥390 samples per class over the 300s soak (≥100 floor for
# Replay rotates round-robin at QPS=8 → ~1.14 QPS per class →
# ≥340 samples per class over the 300s soak (≥100 floor for
# accuracy reduction).
# Replay-range / warm-precompute alignment (issue #46 ε-bound
# bug, fix/quantile-window-alignment): the DDSketch (entry 1) and
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,12 @@ processors:
path: ./processor/countminsketchmergeprocessor
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/processor/countsketchmergeprocessor v0.141.0
path: ./processor/countsketchmergeprocessor
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/processor/kllmergeprocessor v0.141.0
path: ./processor/kllmergeprocessor
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/processor/ddsketchmergeprocessor v0.141.0
path: ./processor/ddsketchmergeprocessor
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/processor/hllmergeprocessor v0.141.0
path: ./processor/hllmergeprocessor
# Paper baseline B1: Serf (ASAP XOR-with-quantization).
- gomod: github.com/open-telemetry/opentelemetry-collector-contrib/processor/serfprocessor v0.141.0
path: ./processor/serfprocessor
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package ddsketchmergeprocessor

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

// Config configures the DDSketch merge processor.
// This processor runs on the gateway side and accumulates per-series
// DDSketch state from `pmetric.MetricTypeDDSketch` data points produced
// by the agent-tier `ddsketchprocessor`. The accumulator is keyed by
// the data point's attribute set; downstream consumers read merged
// state via `GetAccumulator(key)`.
type Config struct {
// MetricName is the metric name to watch for DDSketch payloads.
// Defaults to "ddsketch" if empty. Producers (ddsketchprocessor)
// emit names like "<base>_ddsketch" by default; pipelines should
// override MetricName to match the chosen agent-side suffix.
MetricName string `mapstructure:"metric_name"`
}

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

func (c *Config) Validate() error {
return nil
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package ddsketchmergeprocessor

import (
"context"

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

var typeStr = component.MustNewType("ddsketchmerge")

func NewFactory() processor.Factory {
return processor.NewFactory(
typeStr,
createDefaultConfig,
processor.WithMetrics(createMetricsProcessor, component.StabilityLevelDevelopment),
)
}

func createDefaultConfig() component.Config {
return &Config{}
}

func createMetricsProcessor(
ctx context.Context,
set processor.Settings,
cfg component.Config,
next consumer.Metrics,
) (processor.Metrics, error) {
c := cfg.(*Config)
logger := set.Logger
if logger == nil {
logger = zap.NewNop()
}
p := newProcessor(c, logger, next)
return processorhelper.NewMetrics(ctx, set, cfg, next,
p.processMetrics,
processorhelper.WithStart(p.Start),
processorhelper.WithShutdown(p.Shutdown),
processorhelper.WithCapabilities(p.Capabilities()),
)
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
module github.com/open-telemetry/opentelemetry-collector-contrib/processor/ddsketchmergeprocessor

go 1.24.0

require (
github.com/ProjectASAP/sketchlib-go v0.0.0-20260328221809-b24e56e64e94
go.opentelemetry.io/collector/component v1.47.0
go.opentelemetry.io/collector/consumer v1.47.0
go.opentelemetry.io/collector/pdata v1.47.0
go.opentelemetry.io/collector/processor v1.47.0
go.opentelemetry.io/collector/processor/processorhelper v0.141.0
go.uber.org/zap v1.27.1
)

require (
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853 // indirect
github.com/hashicorp/go-version v1.7.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/klauspost/cpuid/v2 v2.2.10 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect
github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.66.1 // indirect
github.com/prometheus/prometheus v0.307.1 // indirect
github.com/zeebo/xxh3 v1.1.0 // indirect
go.opentelemetry.io/collector/featuregate v1.47.0 // indirect
go.opentelemetry.io/collector/pipeline v1.47.0 // indirect
go.opentelemetry.io/otel v1.38.0 // indirect
go.opentelemetry.io/otel/metric v1.38.0 // indirect
go.opentelemetry.io/otel/trace v1.38.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.yaml.in/yaml/v2 v2.4.3 // indirect
golang.org/x/sys v0.37.0 // indirect
golang.org/x/text v0.30.0 // indirect
google.golang.org/protobuf v1.36.11 // indirect
)

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

replace github.com/ProjectASAP/sketchlib-go => ../../../../sketchlib-go
Loading