Skip to content

Commit 6a3572b

Browse files
committed
opentelemetry: Implement dual Load Balancer delay spans and metrics
1 parent 5e56f38 commit 6a3572b

17 files changed

Lines changed: 213 additions & 26 deletions

File tree

api/src/main/java/io/grpc/LoadBalancer.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -697,7 +697,9 @@ public PickResult copyWithSubchannel(Subchannel subchannel) {
697697
*/
698698
public PickResult copyWithStreamTracerFactory(
699699
@Nullable ClientStreamTracer.Factory streamTracerFactory) {
700-
return new PickResult(subchannel, streamTracerFactory, status, drop, authorityOverride, delayType, delayReason);
700+
return new PickResult(
701+
subchannel, streamTracerFactory, status, drop, authorityOverride, delayType,
702+
delayReason);
701703
}
702704

703705
/**

core/src/main/java/io/grpc/internal/DelayedClientTransport.java

Lines changed: 6 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -158,8 +158,8 @@ public final ClientStream newStream(
158158
synchronized (lock) {
159159
PickerState newerState = pickerState;
160160
if (state == newerState) {
161-
String delayType = determineQueuingDelayType(pickResult, callOptions.isWaitForReady());
162-
String delayReason = determineQueuingDelayReason(pickResult, callOptions.isWaitForReady());
161+
String delayType = determineQueuingDelayType(pickResult);
162+
String delayReason = determineQueuingDelayReason(pickResult);
163163
return createPendingStream(args, tracers, pickResult, delayType, delayReason);
164164
}
165165
state = newerState;
@@ -320,10 +320,8 @@ final void reprocess(@Nullable SubchannelPicker picker) {
320320
}
321321
toRemove.add(stream);
322322
} else { // stay pending
323-
String delayType = determineQueuingDelayType(
324-
pickResult, stream.args.getCallOptions().isWaitForReady());
325-
String delayReason = determineQueuingDelayReason(
326-
pickResult, stream.args.getCallOptions().isWaitForReady());
323+
String delayType = determineQueuingDelayType(pickResult);
324+
String delayReason = determineQueuingDelayReason(pickResult);
327325
stream.updateDelay(delayType, delayReason);
328326
}
329327
}
@@ -366,8 +364,7 @@ public InternalLogId getLogId() {
366364
return logId;
367365
}
368366

369-
private static String determineQueuingDelayType(
370-
@Nullable PickResult pickResult, boolean isWaitForReady) {
367+
private static String determineQueuingDelayType(@Nullable PickResult pickResult) {
371368
if (pickResult == null) {
372369
return "client_channel_init";
373370
}
@@ -383,8 +380,7 @@ private static String determineQueuingDelayType(
383380
return "client_channel_init";
384381
}
385382

386-
private static String determineQueuingDelayReason(
387-
@Nullable PickResult pickResult, boolean isWaitForReady) {
383+
private static String determineQueuingDelayReason(@Nullable PickResult pickResult) {
388384
if (pickResult == null) {
389385
return "client channel: created LB policy.";
390386
}

core/src/main/java/io/grpc/internal/PickFirstLeafLoadBalancer.java

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -167,7 +167,10 @@ public Status acceptResolvedAddresses(ResolvedAddresses resolvedAddresses) {
167167
if (noOldAddrs) {
168168
// Make tests happy; they don't properly assume starting in CONNECTING
169169
rawConnectivityState = CONNECTING;
170-
updateBalancingState(CONNECTING, new FixedResultPicker(PickResult.withNoResult()));
170+
updateBalancingState(
171+
CONNECTING,
172+
new FixedResultPicker(
173+
PickResult.withNoResult("connecting", "pick_first: address list updated")));
171174
}
172175

173176
if (rawConnectivityState == READY) {
@@ -333,7 +336,10 @@ void processSubchannelState(SubchannelData subchannelData, ConnectivityStateInfo
333336

334337
case CONNECTING:
335338
rawConnectivityState = CONNECTING;
336-
updateBalancingState(CONNECTING, new FixedResultPicker(PickResult.withNoResult()));
339+
updateBalancingState(
340+
CONNECTING,
341+
new FixedResultPicker(
342+
PickResult.withNoResult("connecting", "pick_first: attempting to connect")));
337343
break;
338344

339345
case READY:
@@ -653,7 +659,8 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
653659
if (connectionRequested.compareAndSet(false, true)) {
654660
helper.getSynchronizationContext().execute(pickFirstLeafLoadBalancer::requestConnection);
655661
}
656-
return PickResult.withNoResult();
662+
return PickResult.withNoResult(
663+
"connecting", "pick_first: requesting connection");
657664
}
658665
}
659666

core/src/main/java/io/grpc/internal/PickFirstLoadBalancer.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -180,7 +180,8 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
180180
if (connectionRequested.compareAndSet(false, true)) {
181181
helper.getSynchronizationContext().execute(PickFirstLoadBalancer.this::requestConnection);
182182
}
183-
return PickResult.withNoResult();
183+
return PickResult.withNoResult(
184+
"connecting", "pick_first: requesting connection");
184185
}
185186
}
186187

core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -888,7 +888,8 @@ public void streamDelayMetrics_channelFallback_subchannelStateMismatch() {
888888
ClientStreamTracer mockTracer = mock(ClientStreamTracer.class);
889889
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer };
890890

891-
io.grpc.LoadBalancer.Subchannel disconnectedSubchannel = mock(io.grpc.LoadBalancer.Subchannel.class);
891+
io.grpc.LoadBalancer.Subchannel disconnectedSubchannel =
892+
mock(io.grpc.LoadBalancer.Subchannel.class);
892893
when(disconnectedSubchannel.getInternalSubchannel())
893894
.thenReturn(newTransportProvider(null));
894895

grpclb/src/main/java/io/grpc/grpclb/GrpclbState.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,8 @@ final class GrpclbState {
124124
static final RoundRobinEntry BUFFER_ENTRY = new RoundRobinEntry() {
125125
@Override
126126
public PickResult picked(PickSubchannelArgs args) {
127-
return PickResult.withNoResult();
127+
return PickResult.withNoResult(
128+
"connecting", "grpclb: waiting for backend server list");
128129
}
129130

130131
@Override

opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -232,6 +232,16 @@ static OpenTelemetryMetricsResource createMetricInstruments(Meter meter,
232232
.build());
233233
}
234234

235+
if (isMetricEnabled("grpc.client.attempt.delay", enableMetrics, disableDefault)) {
236+
builder.clientAttemptDelayCounter(
237+
meter.histogramBuilder(
238+
"grpc.client.attempt.delay")
239+
.setUnit("s")
240+
.setDescription("Time taken to complete a client call attempt delay")
241+
.setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
242+
.build());
243+
}
244+
235245
if (isMetricEnabled("grpc.client.attempt.sent_total_compressed_message_size", enableMetrics,
236246
disableDefault)) {
237247
builder.clientTotalSentCompressedMessageSizeCounter(

opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -203,6 +203,8 @@ private static final class ClientTracer extends ClientStreamTracer {
203203
volatile String backendService;
204204
long attemptNanos;
205205
Code statusCode;
206+
@Nullable private volatile Stopwatch activeDelayStopwatch;
207+
@Nullable private volatile String activeDelayType;
206208

207209
ClientTracer(CallAttemptsTracerFactory attemptsState, OpenTelemetryMetricsModule module,
208210
StreamInfo info, String target, String fullMethodName,
@@ -216,6 +218,59 @@ private static final class ClientTracer extends ClientStreamTracer {
216218
this.stopwatch = module.stopwatchSupplier.get().start();
217219
}
218220

221+
@Override
222+
public void streamCreated(io.grpc.Attributes transportAtts, Metadata headers) {
223+
delayEnded();
224+
}
225+
226+
@Override
227+
public void delayTypeStarted(String delayType) {
228+
delayEnded();
229+
activeDelayType = delayType;
230+
activeDelayStopwatch = module.stopwatchSupplier.get().start();
231+
}
232+
233+
@Override
234+
public void delayEnded() {
235+
Stopwatch delayStopwatch = activeDelayStopwatch;
236+
String delayType = activeDelayType;
237+
if (delayStopwatch != null && delayType != null) {
238+
delayStopwatch.stop();
239+
long delayNanos = delayStopwatch.elapsed(TimeUnit.NANOSECONDS);
240+
activeDelayStopwatch = null;
241+
activeDelayType = null;
242+
if (module.resource.clientAttemptDelayCounter() != null) {
243+
AttributesBuilder builder = io.opentelemetry.api.common.Attributes.builder()
244+
.put(METHOD_KEY, fullMethodName)
245+
.put(TARGET_KEY, target)
246+
.put("grpc.delay_type", delayType);
247+
if (module.localityEnabled) {
248+
String savedLocality = locality;
249+
if (savedLocality == null) {
250+
savedLocality = "";
251+
}
252+
builder.put(LOCALITY_KEY, savedLocality);
253+
}
254+
if (module.backendServiceEnabled) {
255+
String savedBackendService = backendService;
256+
if (savedBackendService == null) {
257+
savedBackendService = "";
258+
}
259+
builder.put(BACKEND_SERVICE_KEY, savedBackendService);
260+
}
261+
if (module.customLabelEnabled) {
262+
builder.put(
263+
CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
264+
}
265+
for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
266+
plugin.addLabels(builder);
267+
}
268+
module.resource.clientAttemptDelayCounter()
269+
.record(delayNanos * SECONDS_PER_NANO, builder.build(), attemptsState.otelContext);
270+
}
271+
}
272+
}
273+
219274
@Override
220275
public void inboundHeaders(Metadata headers) {
221276
for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
@@ -262,6 +317,7 @@ public void inboundTrailers(Metadata trailers) {
262317

263318
@Override
264319
public void streamClosed(Status status) {
320+
delayEnded();
265321
stopwatch.stop();
266322
attemptNanos = stopwatch.elapsed(TimeUnit.NANOSECONDS);
267323
Deadline deadline = info.getCallOptions().getDeadline();

opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsResource.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,9 @@ abstract class OpenTelemetryMetricsResource {
3535
@Nullable
3636
abstract DoubleHistogram clientAttemptDurationCounter();
3737

38+
@Nullable
39+
abstract DoubleHistogram clientAttemptDelayCounter();
40+
3841
@Nullable
3942
abstract LongHistogram clientTotalSentCompressedMessageSizeCounter();
4043

@@ -79,6 +82,8 @@ abstract static class Builder {
7982

8083
abstract Builder clientAttemptDurationCounter(DoubleHistogram counter);
8184

85+
abstract Builder clientAttemptDelayCounter(DoubleHistogram counter);
86+
8287
abstract Builder clientTotalSentCompressedMessageSizeCounter(LongHistogram counter);
8388

8489
abstract Builder clientTotalReceivedCompressedMessageSizeCounter(

opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -192,6 +192,7 @@ private final class ClientTracer extends ClientStreamTracer {
192192
private final Span parentSpan;
193193
volatile int seqNo;
194194
boolean isPendingStream;
195+
@Nullable private volatile Span activeDelaySpan;
195196

196197
ClientTracer(Span span, Span parentSpan) {
197198
this.span = checkNotNull(span, "span");
@@ -200,6 +201,7 @@ private final class ClientTracer extends ClientStreamTracer {
200201

201202
@Override
202203
public void streamCreated(Attributes transportAtts, Metadata headers) {
204+
delayEnded();
203205
contextPropagators.getTextMapPropagator().inject(Context.current().with(span), headers,
204206
metadataSetter);
205207
if (isPendingStream) {
@@ -212,6 +214,36 @@ public void createPendingStream() {
212214
isPendingStream = true;
213215
}
214216

217+
@Override
218+
public void delayTypeStarted(String delayType) {
219+
if (activeDelaySpan != null) {
220+
activeDelaySpan.end();
221+
}
222+
activeDelaySpan = otelTracer.spanBuilder("Attempt Delay: " + delayType)
223+
.setParent(Context.current().with(span))
224+
.setAttribute("grpc.delay_type", delayType)
225+
.startSpan();
226+
}
227+
228+
@Override
229+
public void delayReasonAttached(String delayReason) {
230+
Span delaySpan = activeDelaySpan;
231+
if (delaySpan != null) {
232+
delaySpan.addEvent(delayReason);
233+
} else {
234+
span.addEvent("delay_reason: " + delayReason);
235+
}
236+
}
237+
238+
@Override
239+
public void delayEnded() {
240+
Span delaySpan = activeDelaySpan;
241+
if (delaySpan != null) {
242+
delaySpan.end();
243+
activeDelaySpan = null;
244+
}
245+
}
246+
215247
@Override
216248
public void outboundMessageSent(
217249
int seqNo, long optionalWireSize, long optionalUncompressedSize) {
@@ -238,6 +270,7 @@ public void inboundUncompressedSize(long bytes) {
238270

239271
@Override
240272
public void streamClosed(io.grpc.Status status) {
273+
delayEnded();
241274
endSpanWithStatus(span, status);
242275
}
243276
}

0 commit comments

Comments
 (0)