Skip to content

Commit c73e0e7

Browse files
committed
chore: Move channel priming out of InstantiatingGrpcTransportProvider
Change-Id: I7214aa3016bd7e7f7f167c64cbaa04134b54a352
1 parent 717bc85 commit c73e0e7

4 files changed

Lines changed: 75 additions & 25 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@
4747
* channel by sending a ReadRow request for a hardcoded, non-existent row key.
4848
*/
4949
@BetaApi("Channel priming is not currently stable and might change in the future")
50-
class BigtableChannelPrimer implements ChannelPrimer {
50+
public class BigtableChannelPrimer implements ChannelPrimer {
5151
private static Logger LOG = Logger.getLogger(BigtableChannelPrimer.class.toString());
5252

5353
static final Metadata.Key<String> REQUEST_PARAMS =

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

Lines changed: 28 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -121,21 +121,38 @@ public static BigtableClientContext create(EnhancedBigtableStubSettings settings
121121
setupCookieHolder(transportProvider);
122122
}
123123

124+
BigtableChannelPrimer channelPrimer =
125+
BigtableChannelPrimer.create(
126+
builder.getProjectId(),
127+
builder.getInstanceId(),
128+
builder.getAppProfileId(),
129+
credentials,
130+
builder.getHeaderProvider().getHeaders());
131+
132+
// // Inject channel priming if enabled
133+
// if (builder.isRefreshingChannel()) {
134+
// transportProvider.setChannelPrimer(
135+
// BigtableChannelPrimer.create(
136+
// builder.getProjectId(),
137+
// builder.getInstanceId(),
138+
// builder.getAppProfileId(),
139+
// credentials,
140+
// builder.getHeaderProvider().getHeaders()));
141+
// }
142+
143+
BigtableTransportChannelProvider btTransportProvider;
144+
124145
// Inject channel priming if enabled
125146
if (builder.isRefreshingChannel()) {
126-
transportProvider.setChannelPrimer(
127-
BigtableChannelPrimer.create(
128-
builder.getProjectId(),
129-
builder.getInstanceId(),
130-
builder.getAppProfileId(),
131-
credentials,
132-
builder.getHeaderProvider().getHeaders()));
147+
btTransportProvider =
148+
BigtableTransportChannelProvider.create(
149+
(InstantiatingGrpcChannelProvider) transportProvider.build(), channelPrimer);
150+
} else {
151+
btTransportProvider =
152+
BigtableTransportChannelProvider.createWithoutChannelPriming(
153+
(InstantiatingGrpcChannelProvider) transportProvider.build());
133154
}
134155

135-
BigtableTransportChannelProvider btTransportProvider =
136-
BigtableTransportChannelProvider.create(
137-
(InstantiatingGrpcChannelProvider) transportProvider.build());
138-
139156
builder.setTransportChannelProvider(btTransportProvider);
140157
}
141158

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

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
import com.google.api.core.InternalApi;
1919
import com.google.api.gax.grpc.ChannelFactory;
20+
import com.google.cloud.bigtable.data.v2.stub.BigtableChannelPrimer;
2021
import com.google.common.annotations.VisibleForTesting;
2122
import com.google.common.base.Preconditions;
2223
import com.google.common.collect.ImmutableList;
@@ -60,6 +61,8 @@ public class BigtableChannelPool extends ManagedChannel {
6061

6162
private final BigtableChannelPoolSettings settings;
6263
private final ChannelFactory channelFactory;
64+
65+
@Nullable private final BigtableChannelPrimer channelPrimer;
6366
private final ScheduledExecutorService executor;
6467

6568
private final Object entryWriteLock = new Object();
@@ -68,9 +71,12 @@ public class BigtableChannelPool extends ManagedChannel {
6871
private final String authority;
6972

7073
public static BigtableChannelPool create(
71-
BigtableChannelPoolSettings settings, ChannelFactory channelFactory) throws IOException {
74+
BigtableChannelPoolSettings settings,
75+
ChannelFactory channelFactory,
76+
BigtableChannelPrimer channelPrimer)
77+
throws IOException {
7278
return new BigtableChannelPool(
73-
settings, channelFactory, Executors.newSingleThreadScheduledExecutor());
79+
settings, channelFactory, channelPrimer, Executors.newSingleThreadScheduledExecutor());
7480
}
7581

7682
/**
@@ -84,15 +90,21 @@ public static BigtableChannelPool create(
8490
BigtableChannelPool(
8591
BigtableChannelPoolSettings settings,
8692
ChannelFactory channelFactory,
93+
BigtableChannelPrimer channelPrimer,
8794
ScheduledExecutorService executor)
8895
throws IOException {
8996
this.settings = settings;
9097
this.channelFactory = channelFactory;
98+
this.channelPrimer = channelPrimer;
9199

92100
ImmutableList.Builder<Entry> initialListBuilder = ImmutableList.builder();
93101

94102
for (int i = 0; i < settings.getInitialChannelCount(); i++) {
95-
initialListBuilder.add(new Entry(channelFactory.createSingleChannel()));
103+
ManagedChannel newChannel = channelFactory.createSingleChannel();
104+
if (channelPrimer != null) {
105+
channelPrimer.primeChannel(newChannel);
106+
}
107+
initialListBuilder.add(new Entry(newChannel));
96108
}
97109

98110
entries.set(initialListBuilder.build());
@@ -316,7 +328,11 @@ private void expand(int desiredSize) {
316328

317329
for (int i = 0; i < desiredSize - localEntries.size(); i++) {
318330
try {
319-
newEntries.add(new Entry(channelFactory.createSingleChannel()));
331+
ManagedChannel newChannel = channelFactory.createSingleChannel();
332+
if (this.channelPrimer != null) {
333+
this.channelPrimer.primeChannel(newChannel);
334+
}
335+
newEntries.add(new Entry(newChannel));
320336
} catch (IOException e) {
321337
LOG.log(Level.WARNING, "Failed to add channel", e);
322338
}
@@ -354,7 +370,11 @@ void refresh() {
354370

355371
for (int i = 0; i < newEntries.size(); i++) {
356372
try {
357-
newEntries.set(i, new Entry(channelFactory.createSingleChannel()));
373+
ManagedChannel newChannel = channelFactory.createSingleChannel();
374+
if (this.channelPrimer != null) {
375+
this.channelPrimer.primeChannel(newChannel);
376+
}
377+
newEntries.set(i, new Entry(newChannel));
358378
} catch (IOException e) {
359379
LOG.log(Level.WARNING, "Failed to refresh channel, leaving old channel", e);
360380
}

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

Lines changed: 21 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -23,12 +23,14 @@
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.BigtableChannelPrimer;
2627
import com.google.common.base.Preconditions;
2728
import io.grpc.ManagedChannel;
2829
import java.io.IOException;
2930
import java.util.Map;
3031
import java.util.concurrent.Executor;
3132
import java.util.concurrent.ScheduledExecutorService;
33+
import javax.annotation.Nullable;
3234

3335
/**
3436
* An instance of TransportChannelProvider that provides a TransportChannel through a supplied
@@ -38,10 +40,13 @@
3840
public final class BigtableTransportChannelProvider implements TransportChannelProvider {
3941

4042
private final InstantiatingGrpcChannelProvider delegate;
43+
@Nullable private final BigtableChannelPrimer channelPrimer;
4144

4245
private BigtableTransportChannelProvider(
43-
InstantiatingGrpcChannelProvider instantiatingGrpcChannelProvider) {
46+
InstantiatingGrpcChannelProvider instantiatingGrpcChannelProvider,
47+
BigtableChannelPrimer channelPrimer) {
4448
delegate = Preconditions.checkNotNull(instantiatingGrpcChannelProvider);
49+
this.channelPrimer = channelPrimer;
4550
}
4651

4752
@Override
@@ -63,7 +68,7 @@ public BigtableTransportChannelProvider withExecutor(ScheduledExecutorService ex
6368
public BigtableTransportChannelProvider withExecutor(Executor executor) {
6469
InstantiatingGrpcChannelProvider newChannelProvider =
6570
(InstantiatingGrpcChannelProvider) delegate.withExecutor(executor);
66-
return new BigtableTransportChannelProvider(newChannelProvider);
71+
return new BigtableTransportChannelProvider(newChannelProvider, channelPrimer);
6772
}
6873

6974
@Override
@@ -75,7 +80,7 @@ public boolean needsHeaders() {
7580
public BigtableTransportChannelProvider withHeaders(Map<String, String> headers) {
7681
InstantiatingGrpcChannelProvider newChannelProvider =
7782
(InstantiatingGrpcChannelProvider) delegate.withHeaders(headers);
78-
return new BigtableTransportChannelProvider(newChannelProvider);
83+
return new BigtableTransportChannelProvider(newChannelProvider, channelPrimer);
7984
}
8085

8186
@Override
@@ -87,7 +92,7 @@ public boolean needsEndpoint() {
8792
public TransportChannelProvider withEndpoint(String endpoint) {
8893
InstantiatingGrpcChannelProvider newChannelProvider =
8994
(InstantiatingGrpcChannelProvider) delegate.withEndpoint(endpoint);
90-
return new BigtableTransportChannelProvider(newChannelProvider);
95+
return new BigtableTransportChannelProvider(newChannelProvider, channelPrimer);
9196
}
9297

9398
@Deprecated
@@ -101,7 +106,7 @@ public boolean acceptsPoolSize() {
101106
public TransportChannelProvider withPoolSize(int size) {
102107
InstantiatingGrpcChannelProvider newChannelProvider =
103108
(InstantiatingGrpcChannelProvider) delegate.withPoolSize(size);
104-
return new BigtableTransportChannelProvider(newChannelProvider);
109+
return new BigtableTransportChannelProvider(newChannelProvider, channelPrimer);
105110
}
106111

107112
/** Expected to only be called once when BigtableClientContext is created */
@@ -130,7 +135,8 @@ public TransportChannel getTransportChannel() throws IOException {
130135
BigtableChannelPoolSettings btPoolSettings =
131136
BigtableChannelPoolSettings.copyFrom(delegate.getChannelPoolSettings());
132137

133-
BigtableChannelPool btChannelPool = BigtableChannelPool.create(btPoolSettings, channelFactory);
138+
BigtableChannelPool btChannelPool =
139+
BigtableChannelPool.create(btPoolSettings, channelFactory, channelPrimer);
134140

135141
return GrpcTransportChannel.create(btChannelPool);
136142
}
@@ -149,12 +155,19 @@ public boolean needsCredentials() {
149155
public TransportChannelProvider withCredentials(Credentials credentials) {
150156
InstantiatingGrpcChannelProvider newChannelProvider =
151157
(InstantiatingGrpcChannelProvider) delegate.withCredentials(credentials);
152-
return new BigtableTransportChannelProvider(newChannelProvider);
158+
return new BigtableTransportChannelProvider(newChannelProvider, channelPrimer);
153159
}
154160

155161
/** Creates a BigtableTransportChannelProvider. */
156162
public static BigtableTransportChannelProvider create(
163+
InstantiatingGrpcChannelProvider instantiatingGrpcChannelProvider,
164+
BigtableChannelPrimer channelPrimer) {
165+
Preconditions.checkNotNull(channelPrimer);
166+
return new BigtableTransportChannelProvider(instantiatingGrpcChannelProvider, channelPrimer);
167+
}
168+
169+
public static BigtableTransportChannelProvider createWithoutChannelPriming(
157170
InstantiatingGrpcChannelProvider instantiatingGrpcChannelProvider) {
158-
return new BigtableTransportChannelProvider(instantiatingGrpcChannelProvider);
171+
return new BigtableTransportChannelProvider(instantiatingGrpcChannelProvider, null);
159172
}
160173
}

0 commit comments

Comments
 (0)