Skip to content
Draft
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
18 changes: 9 additions & 9 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -57,22 +57,22 @@ require (
github.com/stretchr/testify v1.11.1
go.opentelemetry.io/contrib/bridges/prometheus v0.68.0
go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.63.0
go.opentelemetry.io/otel v1.43.0
go.opentelemetry.io/otel v1.44.0
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.12.2
go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.19.0
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.36.0
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.44.0
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.43.0
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.36.0
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.43.0
go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.13.0
go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.36.0
go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.36.0
go.opentelemetry.io/otel/log v0.19.0
go.opentelemetry.io/otel/metric v1.43.0
go.opentelemetry.io/otel/sdk v1.43.0
go.opentelemetry.io/otel/metric v1.44.0
go.opentelemetry.io/otel/sdk v1.44.0
go.opentelemetry.io/otel/sdk/log v0.19.0
go.opentelemetry.io/otel/sdk/metric v1.43.0
go.opentelemetry.io/otel/trace v1.43.0
go.opentelemetry.io/otel/sdk/metric v1.44.0
go.opentelemetry.io/otel/trace v1.44.0
go.uber.org/goleak v1.3.0
go.uber.org/zap v1.27.1
golang.org/x/crypto v0.53.0
Expand All @@ -81,7 +81,7 @@ require (
golang.org/x/time v0.15.0
golang.org/x/tools v0.45.0
gonum.org/v1/gonum v0.17.0
google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa
google.golang.org/grpc v1.82.1
google.golang.org/protobuf v1.36.11
gopkg.in/yaml.v3 v3.0.1
Expand Down Expand Up @@ -110,7 +110,7 @@ require (
github.com/google/flatbuffers v25.2.10+incompatible // indirect
github.com/grafana/pyroscope-go/godeltaprof v0.1.9 // indirect
github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.2 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect
github.com/hako/durafmt v0.0.0-20200710122514-c0fb7b4da026 // indirect
github.com/hashicorp/yamux v0.1.2 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
Expand Down Expand Up @@ -156,7 +156,7 @@ require (
golang.org/x/term v0.44.0 // indirect
golang.org/x/text v0.38.0 // indirect
golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect
)

Expand Down
18 changes: 18 additions & 0 deletions go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions pkg/beholder/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -551,6 +551,8 @@ func newMeterProvider(cfg Config, resource *sdkresource.Resource, auth Auth, cre
for _, p := range cfg.MetricProducers {
readerOpts = append(readerOpts, sdkmetric.WithProducer(p))
}
// NewPeriodicReader reads OTEL_GO_X_METRIC_EXPORT_BATCH_SIZE at
// construction and applies the upstream experimental data-point batching.
mpOpts := append(cfg.metricOptions(),
sdkmetric.WithReader(sdkmetric.NewPeriodicReader(exporter, readerOpts...)),
sdkmetric.WithResource(resource),
Expand Down
1 change: 1 addition & 0 deletions pkg/beholder/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,7 @@ func TestNewClient(t *testing.T) {
t.Run("HTTP endpoint set", func(t *testing.T) {
client, err := beholder.NewClient(beholder.Config{
OtelExporterHTTPEndpoint: "http-endpoint",
InsecureConnection: true,
})
require.NoError(t, err)
assert.NotNil(t, client)
Expand Down
6 changes: 6 additions & 0 deletions pkg/beholder/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,12 @@ type Config struct {
TraceCompressor string

// OTel Metric
// Metric export batching is controlled by the OTel SDK's experimental
// OTEL_GO_X_METRIC_EXPORT_BATCH_SIZE environment variable. Set it to a
// positive integer before the client is created to limit each exporter
// call to that many data points. The value is process-wide and is read
// when each PeriodicReader is constructed; unset, invalid, zero, and
// negative values preserve unbatched export behavior.
MetricReaderInterval time.Duration
MetricRetryConfig *RetryConfig
MetricViews []metric.View
Expand Down
2 changes: 2 additions & 0 deletions pkg/beholder/httpclient.go
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,8 @@ func newHTTPMeterProvider(config Config, resource *sdkresource.Resource, tlsConf

mpOpts := append(config.metricOptions(),
sdkmetric.WithReader(
// NewPeriodicReader reads OTEL_GO_X_METRIC_EXPORT_BATCH_SIZE at
// construction and applies the upstream experimental data-point batching.
sdkmetric.NewPeriodicReader(
exporter,
sdkmetric.WithInterval(config.MetricReaderInterval), // Default is 10s
Expand Down
167 changes: 167 additions & 0 deletions pkg/beholder/metric_export_batch_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
package beholder

import (
"context"
"os"
"sync"
"testing"
"time"

"github.com/stretchr/testify/require"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
)

const metricExportBatchSizeEnv = "OTEL_GO_X_METRIC_EXPORT_BATCH_SIZE"

func TestPeriodicReaderMetricExportBatchSize(t *testing.T) {
testCases := []struct {
name string
envValue *string
wantExporters int
wantMaxPoints int
}{
{
name: "positive value batches exports",
envValue: stringPtr("2"),
wantExporters: 3,
wantMaxPoints: 2,
},
{
name: "invalid value is unbatched",
envValue: stringPtr("invalid"),
wantExporters: 1,
wantMaxPoints: 5,
},
{
name: "zero is unbatched",
envValue: stringPtr("0"),
wantExporters: 1,
wantMaxPoints: 5,
},
{
name: "negative value is unbatched",
envValue: stringPtr("-1"),
wantExporters: 1,
wantMaxPoints: 5,
},
{
name: "unset value is unbatched",
wantExporters: 1,
wantMaxPoints: 5,
},
}

for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
setMetricExportBatchSize(t, tc.envValue)

exporter := &metricBatchRecorder{}
reader := sdkmetric.NewPeriodicReader(exporter, sdkmetric.WithInterval(time.Hour))
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
t.Cleanup(func() { require.NoError(t, provider.Shutdown(context.Background())) })

meter := provider.Meter("beholder/metric-export-batching-test")
counter, err := meter.Int64Counter("metric_export_batching_points")
require.NoError(t, err)
for i := range 5 {
counter.Add(context.Background(), 1, metric.WithAttributes(attribute.Int("point", i)))
}

require.NoError(t, provider.ForceFlush(context.Background()))

batches := exporter.batchesSnapshot()
require.Len(t, batches, tc.wantExporters)
seen := make(map[int]int, 5)
for _, batch := range batches {
require.LessOrEqual(t, len(batch), tc.wantMaxPoints)
for _, point := range batch {
seen[point.attribute]++
}
}
require.Len(t, seen, 5)
for i := range 5 {
require.Equal(t, 1, seen[i])
}
})
}
}

func setMetricExportBatchSize(t *testing.T, value *string) {
t.Helper()
previous, wasSet := os.LookupEnv(metricExportBatchSizeEnv)
if value == nil {
require.NoError(t, os.Unsetenv(metricExportBatchSizeEnv))
} else {
require.NoError(t, os.Setenv(metricExportBatchSizeEnv, *value))
}
t.Cleanup(func() {
if wasSet {
_ = os.Setenv(metricExportBatchSizeEnv, previous)
return
}
_ = os.Unsetenv(metricExportBatchSizeEnv)
})
}

func stringPtr(value string) *string { return &value }

type metricBatchRecorder struct {
mu sync.Mutex
batches [][]metricBatchPoint
}

type metricBatchPoint struct {
attribute int
}

var _ sdkmetric.Exporter = (*metricBatchRecorder)(nil)

func (e *metricBatchRecorder) Temporality(sdkmetric.InstrumentKind) metricdata.Temporality {
return metricdata.CumulativeTemporality
}

func (e *metricBatchRecorder) Aggregation(sdkmetric.InstrumentKind) sdkmetric.Aggregation {
return sdkmetric.AggregationDefault{}
}

func (e *metricBatchRecorder) Export(_ context.Context, metrics *metricdata.ResourceMetrics) error {
var points []metricBatchPoint
for _, scope := range metrics.ScopeMetrics {
for _, metric := range scope.Metrics {
sum, ok := metric.Data.(metricdata.Sum[int64])
if !ok {
continue
}
for _, point := range sum.DataPoints {
value, ok := point.Attributes.Value("point")
if !ok {
continue
}
points = append(points, metricBatchPoint{attribute: int(value.AsInt64())})
}
}
}

e.mu.Lock()
defer e.mu.Unlock()
e.batches = append(e.batches, points)
return nil
}

func (e *metricBatchRecorder) ForceFlush(context.Context) error { return nil }

func (e *metricBatchRecorder) Shutdown(context.Context) error { return nil }

func (e *metricBatchRecorder) batchesSnapshot() [][]metricBatchPoint {
e.mu.Lock()
defer e.mu.Unlock()

batches := make([][]metricBatchPoint, len(e.batches))
for i, batch := range e.batches {
batches[i] = append([]metricBatchPoint(nil), batch...)
}
return batches
}
20 changes: 20 additions & 0 deletions pkg/beholder/metric_export_batching.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
# OTel metric export batching

The OpenTelemetry SDK can split one metric collection into multiple exporter
calls by setting this environment variable before the node constructs its
Beholder meter provider:

```text
OTEL_GO_X_METRIC_EXPORT_BATCH_SIZE=<positive data-point count>
```

This is an experimental, process-wide SDK setting. It is read when each
`PeriodicReader` is constructed, so changing the environment afterward does
not reconfigure an existing reader. The value limits the number of metric
data points per exporter call; it is not a serialized-byte limit. Large
attributes, histograms, or a single oversized data point can still exceed a
collector receive limit.

An unset, invalid, zero, or negative value preserves the default unbatched
behavior. Configure this at deployment time; library code must not mutate the
process environment.
Loading