Skip to content
This repository was archived by the owner on May 8, 2026. It is now read-only.

Commit 6f6e481

Browse files
committed
WIP
1 parent ba09467 commit 6f6e481

8 files changed

Lines changed: 61 additions & 58 deletions

File tree

google-cloud-bigtable/clirr-ignored-differences.xml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,10 @@
5555
<differenceType>8001</differenceType>
5656
<className>com/google/cloud/bigtable/data/v2/stub/metrics/HeaderTracer</className>
5757
</difference>
58+
<difference>
59+
<differenceType>8001</differenceType>
60+
<className>com/google/cloud/bigtable/data/v2/stub/metrics/ErrorCountPerConnectionMetricTracker</className>
61+
</difference>
5862
<!-- InternalApi that was removed -->
5963
<difference>
6064
<differenceType>8001</differenceType>

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/BigtableClientContext.java

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@
2727
import com.google.cloud.bigtable.data.v2.BigtableDataSettings;
2828
import com.google.cloud.bigtable.data.v2.internal.JwtCredentialsWithAudience;
2929
import com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants;
30-
import com.google.cloud.bigtable.data.v2.stub.metrics.ChannelPoolMetricsTracker;
30+
import com.google.cloud.bigtable.data.v2.stub.metrics.ChannelPoolMetricsTracer;
3131
import com.google.cloud.bigtable.data.v2.stub.metrics.CustomOpenTelemetryMetricsProvider;
3232
import com.google.cloud.bigtable.data.v2.stub.metrics.DefaultMetricsProvider;
3333
import com.google.cloud.bigtable.data.v2.stub.metrics.MetricsProvider;
@@ -97,16 +97,16 @@ public static BigtableClientContext create(EnhancedBigtableStubSettings settings
9797
: null;
9898

9999
@Nullable OpenTelemetrySdk internalOtel = null;
100-
@Nullable ChannelPoolMetricsTracker channelPoolMetricsTracker = null;
100+
@Nullable ChannelPoolMetricsTracer channelPoolMetricsTracer = null;
101101
// Internal metrics are scoped to the connections, so we need a mutable transportProvider,
102102
// otherwise there is
103103
// no reason to build the internal OtelProvider
104104
if (transportProvider != null) {
105105
internalOtel =
106106
settings.getInternalMetricsProvider().createOtelProvider(settings, credentials);
107107
if (internalOtel != null) {
108-
channelPoolMetricsTracker =
109-
new ChannelPoolMetricsTracker(
108+
channelPoolMetricsTracer =
109+
new ChannelPoolMetricsTracer(
110110
internalOtel, EnhancedBigtableStub.createBuiltinAttributes(builder.build()));
111111

112112
// Configure grpc metrics
@@ -137,14 +137,14 @@ public static BigtableClientContext create(EnhancedBigtableStubSettings settings
137137
BigtableTransportChannelProvider.create(
138138
(InstantiatingGrpcChannelProvider) transportProvider.build(),
139139
channelPrimer,
140-
channelPoolMetricsTracker);
140+
channelPoolMetricsTracer);
141141

142142
builder.setTransportChannelProvider(btTransportProvider);
143143
}
144144

145145
ClientContext clientContext = ClientContext.create(builder.build());
146-
if (channelPoolMetricsTracker != null) {
147-
channelPoolMetricsTracker.start(clientContext.getExecutor());
146+
if (channelPoolMetricsTracer != null) {
147+
channelPoolMetricsTracer.start(clientContext.getExecutor());
148148
}
149149

150150
return new BigtableClientContext(

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/ChannelPoolMetricsTracker.java renamed to google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/ChannelPoolMetricsTracer.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,9 @@
1515
*/
1616
package com.google.cloud.bigtable.data.v2.stub.metrics;
1717

18-
import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.*;
18+
import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.METER_NAME;
19+
import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.OUTSTANDING_RPCS_PER_CHANNEL_NAME;
20+
import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.PER_CONNECTION_ERROR_COUNT_NAME;
1921

2022
import com.google.api.core.InternalApi;
2123
import com.google.cloud.bigtable.gaxx.grpc.BigtableChannelObserver;
@@ -33,8 +35,8 @@
3335
import javax.annotation.Nullable;
3436

3537
@InternalApi("For internal use only")
36-
public class ChannelPoolMetricsTracker implements Runnable {
37-
private static final Logger logger = Logger.getLogger(ChannelPoolMetricsTracker.class.getName());
38+
public class ChannelPoolMetricsTracer implements Runnable {
39+
private static final Logger logger = Logger.getLogger(ChannelPoolMetricsTracer.class.getName());
3840

3941
private static final int SAMPLING_PERIOD_SECONDS = 60;
4042
private final LongHistogram outstandingRpcsHistogram;
@@ -49,7 +51,7 @@ public class ChannelPoolMetricsTracker implements Runnable {
4951
@Nullable private Attributes unaryAttributes;
5052
@Nullable private Attributes streamingAttributes;
5153

52-
public ChannelPoolMetricsTracker(OpenTelemetry openTelemetry, Attributes commonAttrs) {
54+
public ChannelPoolMetricsTracer(OpenTelemetry openTelemetry, Attributes commonAttrs) {
5355
Meter meter = openTelemetry.getMeter(METER_NAME);
5456
this.commonAttrs = commonAttrs;
5557
this.outstandingRpcsHistogram =

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/gaxx/grpc/BigtableChannelObserver.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,10 +27,10 @@ public interface BigtableChannelObserver {
2727
int getOutstandingStreamingRpcs();
2828

2929
/** Get the current number of errors request count since the last observed period */
30-
long getAndResetErrorCount(); // New method
30+
long getAndResetErrorCount();
3131

3232
/** Get the current number of successful requests since the last observed period */
33-
long getAndResetSuccessCount(); // New method
33+
long getAndResetSuccessCount();
3434

3535
boolean isAltsChannel();
3636
}

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/gaxx/grpc/BigtableChannelPool.java

Lines changed: 23 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -565,17 +565,9 @@ static class Entry implements BigtableChannelObserver {
565565
}
566566

567567
void checkAndSetIsAlts(ClientCall<?, ?> call) {
568-
if (isAltsHolder.get() == null) {
569-
boolean result = false;
570-
// try {
571-
// result = AltsContextUtil.check(call.getAttributes());
572-
// } catch (Exception e) {
573-
// LOG.log(Level.FINE, "Failed to check ALTS status on call start", e);
574-
// // result remains false
575-
// }
576-
// Atomically set only if still null
577-
isAltsHolder.compareAndSet(null, result);
578-
}
568+
// TODO(populate ALTS holder)
569+
boolean result = false;
570+
isAltsHolder.compareAndSet(null, result);
579571
}
580572

581573
ManagedChannel getManagedChannel() {
@@ -603,16 +595,16 @@ int getAndResetMaxOutstanding() {
603595
*/
604596
@VisibleForTesting
605597
boolean retain(boolean isStreaming) {
606-
// abort if the channel is closing
607-
if (shutdownRequested.get()) {
608-
release(isStreaming);
609-
return false;
610-
}
611598
AtomicInteger counter = isStreaming ? outstandingStreamingRpcs : outstandingUnaryRpcs;
612599
AtomicInteger maxCounter =
613600
isStreaming ? maxOutstandingStreamingRpcs : maxOutstandingUnaryRpcs;
614601
int currentOutstanding = counter.incrementAndGet();
615602
maxCounter.accumulateAndGet(currentOutstanding, Math::max);
603+
// abort if the channel is closing
604+
if (shutdownRequested.get()) {
605+
release(isStreaming);
606+
return false;
607+
}
616608
return true;
617609
}
618610

@@ -621,8 +613,10 @@ boolean retain(boolean isStreaming) {
621613
* previously requested, this method will shutdown the channel if its the last outstanding RPC.
622614
*/
623615
void release(boolean isStreaming) {
624-
AtomicInteger counter = isStreaming ? outstandingStreamingRpcs : outstandingUnaryRpcs;
625-
int newCount = counter.decrementAndGet();
616+
int newCount =
617+
isStreaming
618+
? outstandingStreamingRpcs.decrementAndGet()
619+
: outstandingUnaryRpcs.decrementAndGet();
626620
if (newCount < 0) {
627621
LOG.log(Level.WARNING, "Bug! Reference count is negative (" + newCount + ")!");
628622
}
@@ -659,14 +653,6 @@ public int getOutstandingUnaryRpcs() {
659653
return outstandingUnaryRpcs.get();
660654
}
661655

662-
void incrementErrorCount() {
663-
errorCount.incrementAndGet();
664-
}
665-
666-
void incrementSuccessCount() {
667-
successCount.incrementAndGet();
668-
}
669-
670656
@Override
671657
public int getOutstandingStreamingRpcs() {
672658
return outstandingStreamingRpcs.get();
@@ -689,6 +675,15 @@ public boolean isAltsChannel() {
689675
Boolean val = isAltsHolder.get();
690676
return val != null && val;
691677
}
678+
679+
680+
void incrementErrorCount() {
681+
errorCount.incrementAndGet();
682+
}
683+
684+
void incrementSuccessCount() {
685+
successCount.incrementAndGet();
686+
}
692687
}
693688

694689
/** Thin wrapper to ensure that new calls are properly reference counted. */
@@ -707,7 +702,8 @@ public String authority() {
707702
@Override
708703
public <RequestT, ResponseT> ClientCall<RequestT, ResponseT> newCall(
709704
MethodDescriptor<RequestT, ResponseT> methodDescriptor, CallOptions callOptions) {
710-
boolean isStreaming = methodDescriptor.getType() != MethodDescriptor.MethodType.UNARY;
705+
boolean isStreaming =
706+
methodDescriptor.getType() == MethodDescriptor.MethodType.SERVER_STREAMING;
711707
Entry entry = getRetainedEntry(index, isStreaming);
712708
return new ReleasingClientCall<>(
713709
entry.channel.newCall(methodDescriptor, callOptions), entry, isStreaming);

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/gaxx/grpc/BigtableTransportChannelProvider.java

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
import com.google.api.gax.rpc.TransportChannel;
2424
import com.google.api.gax.rpc.TransportChannelProvider;
2525
import com.google.auth.Credentials;
26-
import com.google.cloud.bigtable.data.v2.stub.metrics.ChannelPoolMetricsTracker;
26+
import com.google.cloud.bigtable.data.v2.stub.metrics.ChannelPoolMetricsTracer;
2727
import com.google.common.base.Preconditions;
2828
import io.grpc.ManagedChannel;
2929
import java.io.IOException;
@@ -41,15 +41,15 @@ public final class BigtableTransportChannelProvider implements TransportChannelP
4141

4242
private final InstantiatingGrpcChannelProvider delegate;
4343
private final ChannelPrimer channelPrimer;
44-
@Nullable private final ChannelPoolMetricsTracker channelPoolMetricsTracker;
44+
@Nullable private final ChannelPoolMetricsTracer channelPoolMetricsTracer;
4545

4646
private BigtableTransportChannelProvider(
4747
InstantiatingGrpcChannelProvider instantiatingGrpcChannelProvider,
4848
ChannelPrimer channelPrimer,
49-
ChannelPoolMetricsTracker channelPoolMetricsTracker) {
49+
ChannelPoolMetricsTracer channelPoolMetricsTracer) {
5050
delegate = Preconditions.checkNotNull(instantiatingGrpcChannelProvider);
5151
this.channelPrimer = channelPrimer;
52-
this.channelPoolMetricsTracker = channelPoolMetricsTracker;
52+
this.channelPoolMetricsTracer = channelPoolMetricsTracer;
5353
}
5454

5555
@Override
@@ -72,7 +72,7 @@ public BigtableTransportChannelProvider withExecutor(Executor executor) {
7272
InstantiatingGrpcChannelProvider newChannelProvider =
7373
(InstantiatingGrpcChannelProvider) delegate.withExecutor(executor);
7474
return new BigtableTransportChannelProvider(
75-
newChannelProvider, channelPrimer, channelPoolMetricsTracker);
75+
newChannelProvider, channelPrimer, channelPoolMetricsTracer);
7676
}
7777

7878
@Override
@@ -85,7 +85,7 @@ public BigtableTransportChannelProvider withHeaders(Map<String, String> headers)
8585
InstantiatingGrpcChannelProvider newChannelProvider =
8686
(InstantiatingGrpcChannelProvider) delegate.withHeaders(headers);
8787
return new BigtableTransportChannelProvider(
88-
newChannelProvider, channelPrimer, channelPoolMetricsTracker);
88+
newChannelProvider, channelPrimer, channelPoolMetricsTracer);
8989
}
9090

9191
@Override
@@ -98,7 +98,7 @@ public TransportChannelProvider withEndpoint(String endpoint) {
9898
InstantiatingGrpcChannelProvider newChannelProvider =
9999
(InstantiatingGrpcChannelProvider) delegate.withEndpoint(endpoint);
100100
return new BigtableTransportChannelProvider(
101-
newChannelProvider, channelPrimer, channelPoolMetricsTracker);
101+
newChannelProvider, channelPrimer, channelPoolMetricsTracer);
102102
}
103103

104104
@Deprecated
@@ -113,7 +113,7 @@ public TransportChannelProvider withPoolSize(int size) {
113113
InstantiatingGrpcChannelProvider newChannelProvider =
114114
(InstantiatingGrpcChannelProvider) delegate.withPoolSize(size);
115115
return new BigtableTransportChannelProvider(
116-
newChannelProvider, channelPrimer, channelPoolMetricsTracker);
116+
newChannelProvider, channelPrimer, channelPoolMetricsTracer);
117117
}
118118

119119
/** Expected to only be called once when BigtableClientContext is created */
@@ -145,9 +145,9 @@ public TransportChannel getTransportChannel() throws IOException {
145145
BigtableChannelPool btChannelPool =
146146
BigtableChannelPool.create(btPoolSettings, channelFactory, channelPrimer);
147147

148-
if (channelPoolMetricsTracker != null) {
149-
channelPoolMetricsTracker.registerChannelInsightsProvider(btChannelPool::getChannelInfos);
150-
channelPoolMetricsTracker.registerLoadBalancingStrategy(
148+
if (channelPoolMetricsTracer != null) {
149+
channelPoolMetricsTracer.registerChannelInsightsProvider(btChannelPool::getChannelInfos);
150+
channelPoolMetricsTracer.registerLoadBalancingStrategy(
151151
btPoolSettings.getLoadBalancingStrategy().name());
152152
}
153153

@@ -169,14 +169,14 @@ public TransportChannelProvider withCredentials(Credentials credentials) {
169169
InstantiatingGrpcChannelProvider newChannelProvider =
170170
(InstantiatingGrpcChannelProvider) delegate.withCredentials(credentials);
171171
return new BigtableTransportChannelProvider(
172-
newChannelProvider, channelPrimer, channelPoolMetricsTracker);
172+
newChannelProvider, channelPrimer, channelPoolMetricsTracer);
173173
}
174174

175175
/** Creates a BigtableTransportChannelProvider. */
176176
public static BigtableTransportChannelProvider create(
177177
InstantiatingGrpcChannelProvider instantiatingGrpcChannelProvider,
178178
ChannelPrimer channelPrimer,
179-
ChannelPoolMetricsTracker outstandingRpcsMetricTracke) {
179+
ChannelPoolMetricsTracer outstandingRpcsMetricTracke) {
180180
return new BigtableTransportChannelProvider(
181181
instantiatingGrpcChannelProvider, channelPrimer, outstandingRpcsMetricTracke);
182182
}

google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/ChannelPoolMetricsTrackerTest.java renamed to google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/ChannelPoolMetricsTracerTest.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -51,15 +51,15 @@
5151
import org.mockito.stubbing.Answer;
5252

5353
@RunWith(JUnit4.class)
54-
public class ChannelPoolMetricsTrackerTest {
54+
public class ChannelPoolMetricsTracerTest {
5555

5656
@Rule public final MockitoRule mockito = MockitoJUnit.rule();
5757

5858
private InMemoryMetricReader metricReader;
5959
@Mock private ScheduledExecutorService mockScheduler;
6060
private ArgumentCaptor<Runnable> runnableCaptor;
6161

62-
private ChannelPoolMetricsTracker tracker;
62+
private ChannelPoolMetricsTracer tracker;
6363
private Attributes baseAttributes;
6464

6565
@Mock private BigtableChannelPoolObserver mockInsightsProvider;
@@ -76,7 +76,7 @@ public void setUp() {
7676

7777
baseAttributes = Attributes.builder().build();
7878

79-
tracker = new ChannelPoolMetricsTracker(openTelemetry, baseAttributes);
79+
tracker = new ChannelPoolMetricsTracer(openTelemetry, baseAttributes);
8080

8181
runnableCaptor = ArgumentCaptor.forClass(Runnable.class);
8282
// Configure mockScheduler to capture the runnable when tracker.start() is called
@@ -99,12 +99,12 @@ public void setUp() {
9999
when(mockInsight2.isAltsChannel()).thenReturn(false);
100100
}
101101

102-
/** Helper to run the captured ChannelPoolMetricsTracker task. */
102+
/** Helper to run the captured ChannelPoolMetricsTracer task. */
103103
void runTrackerTask() {
104104
List<Runnable> capturedRunnables = runnableCaptor.getAllValues();
105105
assertThat(capturedRunnables).hasSize(1); // Expect only one task scheduled
106106
Runnable trackerRunnable = capturedRunnables.get(0);
107-
assertThat(trackerRunnable).isInstanceOf(ChannelPoolMetricsTracker.class);
107+
assertThat(trackerRunnable).isInstanceOf(ChannelPoolMetricsTracer.class);
108108
trackerRunnable.run();
109109
}
110110

google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/gaxx/grpc/BigtableChannelPoolTest.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -231,5 +231,6 @@ public void testMixedRpcs() {
231231
assertThat(entry.getOutstandingStreamingRpcs()).isEqualTo(0);
232232
assertThat(entry.getAndResetSuccessCount()).isEqualTo(0);
233233
assertThat(entry.getAndResetErrorCount()).isEqualTo(1); // The last failure
234+
assertThat(entry.totalOutstandingRpcs()).isEqualTo(0);
234235
}
235236
}

0 commit comments

Comments
 (0)