You've already forked opentelemetry-go
mirror of
https://github.com/open-telemetry/opentelemetry-go.git
synced 2026-06-19 21:45:50 +02:00
Disable exemplar reservoir for asynchronous instruments by default (#8286)
Fixes https://github.com/open-telemetry/opentelemetry-go/issues/8285 Asynchronous instruments do not accept context, so the default TraceBased exemplar filter can never record an exemplar. Use the DropReservoir when the TraceBased exemplar filter is used, similar to https://github.com/open-telemetry/opentelemetry-go/pull/8211. ### Benchmarks Since callbacks are made as part of collect, there is a significant baseline overhead from collection included in the benchmark. ``` goos: linux goarch: amd64 pkg: go.opentelemetry.io/otel/sdk/metric cpu: AMD EPYC 7B12 │ main.txt │ new.txt │ │ sec/op │ sec/op vs base │ AsyncMeasureNewAttributeSet/AlwaysOn-24 3.662µ ± 13% 3.735µ ± 8% ~ (p=0.699 n=6) AsyncMeasureNewAttributeSet/TraceBased-24 3.586µ ± 10% 1.195µ ± 5% -66.69% (p=0.002 n=6) geomean 3.623µ 2.112µ -41.70% │ main.txt │ new.txt │ │ B/op │ B/op vs base │ AsyncMeasureNewAttributeSet/AlwaysOn-24 3.344Ki ± 0% 3.343Ki ± 0% ~ (p=0.571 n=6) AsyncMeasureNewAttributeSet/TraceBased-24 3424.5 ± 0% 592.0 ± 0% -82.71% (p=0.002 n=6) geomean 3.344Ki 1.390Ki -58.43% │ main.txt │ new.txt │ │ allocs/op │ allocs/op vs base │ AsyncMeasureNewAttributeSet/AlwaysOn-24 16.00 ± 0% 16.00 ± 0% ~ (p=1.000 n=6) ¹ AsyncMeasureNewAttributeSet/TraceBased-24 16.00 ± 0% 12.00 ± 0% -25.00% (p=0.002 n=6) geomean 16.00 13.86 -13.40% ¹ all samples are equal ```
This commit is contained in:
@@ -55,6 +55,7 @@ This project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.htm
|
||||
- `go.opentelemetry.io/otel/sdk/log` now unwraps error chains created with `fmt.Errorf` when deriving the `error.type` attribute from errors on log records. (#8133)
|
||||
- `Set.MarshalLog` method in `go.opentelemetry.io/otel/attribute` now uses `Value.String` formatting following the [OpenTelemetry AnyValue representation for non-OTLP protocols](https://opentelemetry.io/docs/specs/otel/common/#anyvalue). (#8169)
|
||||
- Optimize `go.opentelemetry.io/otel/sdk/metric` to return a drop reservoir and short-circuit `Offer` calls to the exemplar reservoir when `exemplar.AlwaysOffFilter` is configured. (#8211) (#8267)
|
||||
- Optimize `go.opentelemetry.io/otel/sdk/metric` to return a drop reservoir for asynchronous instruments when `exemplar.TraceBasedFilter` is configured. (#8286)
|
||||
|
||||
### Deprecated
|
||||
|
||||
|
||||
@@ -897,3 +897,51 @@ func BenchmarkMeasureNewAttributeSet(b *testing.B) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkAsyncMeasureNewAttributeSet(b *testing.B) {
|
||||
ctx := b.Context()
|
||||
|
||||
for _, filterName := range []string{"AlwaysOn", "TraceBased"} {
|
||||
var filter exemplar.Filter
|
||||
if filterName == "AlwaysOn" {
|
||||
filter = exemplar.AlwaysOnFilter
|
||||
} else {
|
||||
filter = exemplar.TraceBasedFilter
|
||||
}
|
||||
|
||||
b.Run(filterName, func(b *testing.B) {
|
||||
var rdr Reader
|
||||
var meter metric.Meter
|
||||
var count int
|
||||
out := new(metricdata.ResourceMetrics)
|
||||
|
||||
b.ReportAllocs()
|
||||
b.ResetTimer()
|
||||
|
||||
for n := 0; n < b.N; n++ {
|
||||
if n%10000 == 0 {
|
||||
b.StopTimer()
|
||||
rdr = NewManualReader()
|
||||
provider := NewMeterProvider(
|
||||
WithReader(rdr),
|
||||
WithExemplarFilter(filter),
|
||||
)
|
||||
meter = provider.Meter("BenchmarkAsyncMeasureNewAttributeSet")
|
||||
|
||||
_, err := meter.Int64ObservableCounter(
|
||||
"int64-observable-counter",
|
||||
metric.WithInt64Callback(func(_ context.Context, obs metric.Int64Observer) error {
|
||||
obs.Observe(1, metric.WithAttributes(attribute.Int("id", count)))
|
||||
return nil
|
||||
}),
|
||||
)
|
||||
assert.NoError(b, err)
|
||||
b.StartTimer()
|
||||
}
|
||||
|
||||
_ = rdr.Collect(ctx, out)
|
||||
count++
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,12 +20,19 @@ type ExemplarReservoirProviderSelector func(Aggregation) exemplar.ReservoirProvi
|
||||
// reservoirFunc returns the appropriately configured exemplar reservoir
|
||||
// creation func based on the passed InstrumentKind and filter configuration.
|
||||
func reservoirFunc[N int64 | float64](
|
||||
kind InstrumentKind,
|
||||
provider exemplar.ReservoirProvider,
|
||||
filter exemplar.Filter,
|
||||
) func(attribute.Set) aggregate.FilteredExemplarReservoir[N] {
|
||||
if reflect.ValueOf(filter).Pointer() == reflect.ValueOf(exemplar.AlwaysOffFilter).Pointer() {
|
||||
return aggregate.DropReservoir[N]
|
||||
}
|
||||
if (kind == InstrumentKindObservableCounter || kind == InstrumentKindObservableUpDownCounter || kind == InstrumentKindObservableGauge) &&
|
||||
reflect.ValueOf(filter).Pointer() == reflect.ValueOf(exemplar.TraceBasedFilter).Pointer() {
|
||||
// Asynchronous instruments do not accept context, so TraceBasedFilter
|
||||
// will never record any exemplars.
|
||||
return aggregate.DropReservoir[N]
|
||||
}
|
||||
return func(attrs attribute.Set) aggregate.FilteredExemplarReservoir[N] {
|
||||
return aggregate.NewFilteredExemplarReservoir[N](filter, provider(attrs))
|
||||
}
|
||||
|
||||
+60
-13
@@ -62,22 +62,69 @@ func TestFixedSizeExemplarConcurrentSafe(t *testing.T) {
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
func TestReservoirFuncAlwaysOff(t *testing.T) {
|
||||
var invoked bool
|
||||
provider := func(attribute.Set) exemplar.Reservoir {
|
||||
invoked = true
|
||||
return nil
|
||||
func TestReservoirFunc(t *testing.T) {
|
||||
type testCase struct {
|
||||
name string
|
||||
kind InstrumentKind
|
||||
filter exemplar.Filter
|
||||
expectDrop bool
|
||||
}
|
||||
|
||||
f := reservoirFunc[int64](provider, exemplar.AlwaysOffFilter)
|
||||
_ = f(*attribute.EmptySet())
|
||||
testCases := []testCase{
|
||||
{
|
||||
name: "AlwaysOff",
|
||||
kind: InstrumentKindCounter,
|
||||
filter: exemplar.AlwaysOffFilter,
|
||||
expectDrop: true,
|
||||
},
|
||||
{
|
||||
name: "AlwaysOn",
|
||||
kind: InstrumentKindCounter,
|
||||
filter: exemplar.AlwaysOnFilter,
|
||||
expectDrop: false,
|
||||
},
|
||||
{
|
||||
name: "TraceBasedSync",
|
||||
kind: InstrumentKindCounter,
|
||||
filter: exemplar.TraceBasedFilter,
|
||||
expectDrop: false,
|
||||
},
|
||||
{
|
||||
name: "TraceBasedAsyncCounter",
|
||||
kind: InstrumentKindObservableCounter,
|
||||
filter: exemplar.TraceBasedFilter,
|
||||
expectDrop: true,
|
||||
},
|
||||
{
|
||||
name: "TraceBasedAsyncUpDownCounter",
|
||||
kind: InstrumentKindObservableUpDownCounter,
|
||||
filter: exemplar.TraceBasedFilter,
|
||||
expectDrop: true,
|
||||
},
|
||||
{
|
||||
name: "TraceBasedAsyncGauge",
|
||||
kind: InstrumentKindObservableGauge,
|
||||
filter: exemplar.TraceBasedFilter,
|
||||
expectDrop: true,
|
||||
},
|
||||
}
|
||||
|
||||
require.False(t, invoked, "ReservoirProvider should not be invoked when AlwaysOffFilter is used")
|
||||
for _, tc := range testCases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
var invoked bool
|
||||
provider := func(attribute.Set) exemplar.Reservoir {
|
||||
invoked = true
|
||||
return nil
|
||||
}
|
||||
|
||||
// Test non-off filter
|
||||
invoked = false
|
||||
f = reservoirFunc[int64](provider, exemplar.AlwaysOnFilter)
|
||||
_ = f(*attribute.EmptySet())
|
||||
f := reservoirFunc[int64](tc.kind, provider, tc.filter)
|
||||
_ = f(*attribute.EmptySet())
|
||||
|
||||
require.True(t, invoked, "ReservoirProvider should be invoked when AlwaysOnFilter is used")
|
||||
if tc.expectDrop {
|
||||
require.False(t, invoked, "ReservoirProvider should not be invoked")
|
||||
} else {
|
||||
require.True(t, invoked, "ReservoirProvider should be invoked")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -399,6 +399,7 @@ func (i *inserter[N]) cachedAggregator(
|
||||
b := aggregate.Builder[N]{
|
||||
Temporality: i.pipeline.reader.temporality(kind),
|
||||
ReservoirFunc: reservoirFunc[N](
|
||||
kind,
|
||||
stream.ExemplarReservoirProviderSelector(stream.Aggregation),
|
||||
i.pipeline.exemplarFilter,
|
||||
),
|
||||
|
||||
Reference in New Issue
Block a user