-
Notifications
You must be signed in to change notification settings - Fork 29
INFOPLAT-13349: feat(beholder): track bytes per export via custom gRPC stats handler #2251
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
kirqz23
wants to merge
7
commits into
main
Choose a base branch
from
infoplat-13349-metered-exporter
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+511
−10
Open
Changes from all commits
Commits
Show all changes
7 commits
Select commit
Hold shift + click to select a range
698806f
feat(beholder): track log export bytes per node via gRPC stats handler
kirqz23 ebe5d3c
feat(beholder): use shared metered exporter for both metrics and logs
kirqz23 ab408fb
test(beholder): add units tests for metered exporter
kirqz23 ef5f0db
Merge branch 'main' into infoplat-13349-metered-exporter
kirqz23 257bcdf
Potential fix for pull request finding
kirqz23 13b1715
chore: rename sizeCaptureHandler to exportSizeHandler
kirqz23 1440cca
fix: malformed test function names
kirqz23 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,127 @@ | ||
| package beholder | ||
|
|
||
| import ( | ||
| "context" | ||
| "sync/atomic" | ||
|
|
||
| "go.opentelemetry.io/otel/attribute" | ||
| otelmetric "go.opentelemetry.io/otel/metric" | ||
| sdklog "go.opentelemetry.io/otel/sdk/log" | ||
| sdkmetric "go.opentelemetry.io/otel/sdk/metric" | ||
| "go.opentelemetry.io/otel/sdk/metric/metricdata" | ||
| "google.golang.org/grpc/stats" | ||
| ) | ||
|
|
||
| // exportSizeKey is the context key under which a metered exporter stashes a | ||
| // per-export byte holder for exportSizeHandler to fill in. | ||
| type exportSizeKey struct{} | ||
|
|
||
| // exportSizeHandler is a minimal, stateless gRPC stats.Handler that records the | ||
| // uncompressed proto size of each outbound message. | ||
| type exportSizeHandler struct{} | ||
|
|
||
| func (exportSizeHandler) TagConn(ctx context.Context, _ *stats.ConnTagInfo) context.Context { | ||
| return ctx | ||
| } | ||
|
|
||
| func (exportSizeHandler) HandleConn(context.Context, stats.ConnStats) {} | ||
|
|
||
| func (exportSizeHandler) TagRPC(ctx context.Context, _ *stats.RPCTagInfo) context.Context { | ||
| return ctx | ||
| } | ||
|
|
||
| // HandleRPC fires on every gRPC stats event. On OutPayload it stores the | ||
| // uncompressed message length, the same field otelgrpc used for | ||
| // rpc.client.request.size | ||
| func (exportSizeHandler) HandleRPC(ctx context.Context, rs stats.RPCStats) { | ||
| op, ok := rs.(*stats.OutPayload) | ||
| if !ok { | ||
| return | ||
| } | ||
| if holder, ok := ctx.Value(exportSizeKey{}).(*atomic.Int64); ok { | ||
| holder.Store(int64(op.Length)) | ||
| } | ||
| } | ||
|
|
||
| const exportBytesMetric = "beholder.export.bytes" | ||
|
kirqz23 marked this conversation as resolved.
|
||
|
|
||
| // newExportBytesCounter creates the counter shared by all metered exporters. | ||
| func newExportBytesCounter(meter otelmetric.Meter) (otelmetric.Int64Counter, error) { | ||
| return meter.Int64Counter( | ||
| exportBytesMetric, | ||
| otelmetric.WithDescription("Uncompressed OTLP proto size in bytes of each successful export batch, by signal."), | ||
| otelmetric.WithUnit("By"), | ||
| ) | ||
| } | ||
|
|
||
| func exportAttrs(signal, csaPublicKeyHex string) otelmetric.MeasurementOption { | ||
| return otelmetric.WithAttributes( | ||
| attribute.String("otel_signal", signal), | ||
| attribute.String("csa_public_key", csaPublicKeyHex), | ||
| ) | ||
| } | ||
|
|
||
| // meteredExporter holds the shared metering logic: run an export with a per-call | ||
| // size holder in the context, then record the captured OutPayload size on | ||
| // success. | ||
| type meteredExporter struct { | ||
| counter otelmetric.Int64Counter | ||
| attrs otelmetric.MeasurementOption | ||
| } | ||
|
|
||
| func (m meteredExporter) record(ctx context.Context, export func(context.Context) error) error { | ||
| var size atomic.Int64 | ||
| err := export(context.WithValue(ctx, exportSizeKey{}, &size)) | ||
| if err == nil { | ||
| m.counter.Add(ctx, size.Load(), m.attrs) | ||
| } | ||
| return err | ||
| } | ||
|
|
||
| // meteredLogExporter wraps an sdklog.Exporter and records each export batch's | ||
| // uncompressed proto size. It sits above the otlploggrpc retry loop, so Export is | ||
| // called once per logical batch and bytes are counted only on success. | ||
| type meteredLogExporter struct { | ||
| meteredExporter | ||
| inner sdklog.Exporter | ||
| } | ||
|
|
||
| func newMeteredLogExporter(inner sdklog.Exporter, counter otelmetric.Int64Counter, csaPublicKeyHex string) *meteredLogExporter { | ||
| return &meteredLogExporter{ | ||
| meteredExporter: meteredExporter{counter: counter, attrs: exportAttrs("logs", csaPublicKeyHex)}, | ||
| inner: inner, | ||
| } | ||
| } | ||
|
|
||
| func (e *meteredLogExporter) Export(ctx context.Context, records []sdklog.Record) error { | ||
| return e.record(ctx, func(c context.Context) error { return e.inner.Export(c, records) }) | ||
| } | ||
|
|
||
| func (e *meteredLogExporter) Shutdown(ctx context.Context) error { return e.inner.Shutdown(ctx) } | ||
|
|
||
| func (e *meteredLogExporter) ForceFlush(ctx context.Context) error { return e.inner.ForceFlush(ctx) } | ||
|
|
||
| // meteredMetricExporter wraps an sdkmetric.Exporter. | ||
| // It is created by the MeterProvider and has no access to the counter until | ||
| // the MeterProvider exits and calls attachCounter. | ||
| type meteredMetricExporter struct { | ||
| sdkmetric.Exporter | ||
| base atomic.Pointer[meteredExporter] | ||
| } | ||
|
|
||
| func newMeteredMetricExporter(inner sdkmetric.Exporter) *meteredMetricExporter { | ||
| return &meteredMetricExporter{Exporter: inner} | ||
| } | ||
|
|
||
| // attachCounter wires the size counter once the MeterProvider exits. | ||
| func (e *meteredMetricExporter) attachCounter(counter otelmetric.Int64Counter, csaPublicKeyHex string) { | ||
| e.base.Store(&meteredExporter{counter: counter, attrs: exportAttrs("metrics", csaPublicKeyHex)}) | ||
| } | ||
|
|
||
| func (e *meteredMetricExporter) Export(ctx context.Context, rm *metricdata.ResourceMetrics) error { | ||
| base := e.base.Load() | ||
| if base == nil { | ||
| return e.Exporter.Export(ctx, rm) | ||
| } | ||
| return base.record(ctx, func(c context.Context) error { return e.Exporter.Export(c, rm) }) | ||
| } | ||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
gRPC builds
istats.NewCombinedHandler(...)from the registered slice of multiple handlers, that is exactly the delegating handler proposed by the Copilot. This proposal is exactly what gRPC already builds for us internally. The downside of having multiple stats handlers however is that gRPC fires every handler on every event, in registration order. So for oneOutPayloadevent there will be twoHandleRPC(...)calls which might be an overhead. We might consider either keeping bothstatsHandlersif we need them, or droppinggrpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...)).This old handler:
rpc.client.*metrics (which part of them are going to be removed),There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
So, we have two approaches here:
exportSizeHandlershouldn't add much to the performance.grpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...))taking into account that we loose 3 points mentioned in the above comment, howeverrpc.client.*metrics are totally deprecated anyway expect duration, which can be added to our custom handler alongsidebeholder.export.bytes, e.g. sth likebeholder.export.durationcc @pkcll