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

Commit a089244

Browse files
committed
more refactoring
1 parent 72cafc8 commit a089244

1 file changed

Lines changed: 49 additions & 26 deletions

File tree

google-cloud-spanner/src/main/java/com/google/cloud/spanner/spi/v1/KeyAwareChannel.java

Lines changed: 49 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,9 @@
4545
/**
4646
* ManagedChannel that routes eligible requests using location-aware routing hints.
4747
*
48-
* <p>Routing is applied only to streaming read and streaming query methods.
48+
* <p>Routing hints are applied to streaming read/query and unary ExecuteSql. Commit/Rollback use
49+
* transaction affinity when available. BeginTransaction is routed only when a mutation key is
50+
* provided.
4951
*/
5052
@InternalApi
5153
final class KeyAwareChannel extends ManagedChannel {
@@ -80,11 +82,6 @@ private KeyAwareChannel(
8082
this.authority = this.defaultChannel.authority();
8183
}
8284

83-
static KeyAwareChannel create(InstantiatingGrpcChannelProvider channelProvider)
84-
throws IOException {
85-
return new KeyAwareChannel(channelProvider, null);
86-
}
87-
8885
static KeyAwareChannel create(
8986
InstantiatingGrpcChannelProvider channelProvider,
9087
@Nullable ChannelEndpointCacheFactory endpointCacheFactory)
@@ -281,29 +278,15 @@ public void sendMessage(RequestT message) {
281278

282279
if (message instanceof ReadRequest) {
283280
ReadRequest.Builder reqBuilder = ((ReadRequest) message).toBuilder();
284-
String databaseId = parentChannel.extractDatabaseIdFromSession(reqBuilder.getSession());
285-
ByteString transactionId = transactionIdFromSelector(reqBuilder.getTransaction());
286-
endpoint = parentChannel.affinityEndpoint(transactionId);
287-
if (databaseId != null) {
288-
finder = parentChannel.getOrCreateChannelFinder(databaseId);
289-
ChannelEndpoint routed = finder.findServer(reqBuilder);
290-
if (endpoint == null) {
291-
endpoint = routed;
292-
}
293-
}
281+
RoutingDecision routing = routeFromRequest(reqBuilder);
282+
finder = routing.finder;
283+
endpoint = routing.endpoint;
294284
message = (RequestT) reqBuilder.build();
295285
} else if (message instanceof ExecuteSqlRequest) {
296286
ExecuteSqlRequest.Builder reqBuilder = ((ExecuteSqlRequest) message).toBuilder();
297-
String databaseId = parentChannel.extractDatabaseIdFromSession(reqBuilder.getSession());
298-
ByteString transactionId = transactionIdFromSelector(reqBuilder.getTransaction());
299-
endpoint = parentChannel.affinityEndpoint(transactionId);
300-
if (databaseId != null) {
301-
finder = parentChannel.getOrCreateChannelFinder(databaseId);
302-
ChannelEndpoint routed = finder.findServer(reqBuilder);
303-
if (endpoint == null) {
304-
endpoint = routed;
305-
}
306-
}
287+
RoutingDecision routing = routeFromRequest(reqBuilder);
288+
finder = routing.finder;
289+
endpoint = routing.endpoint;
307290
message = (RequestT) reqBuilder.build();
308291
} else if (message instanceof BeginTransactionRequest) {
309292
BeginTransactionRequest.Builder reqBuilder =
@@ -373,6 +356,46 @@ void maybeRecordAffinity(ByteString transactionId) {
373356
void maybeClearAffinity() {
374357
parentChannel.clearAffinity(transactionIdToClear);
375358
}
359+
360+
private RoutingDecision routeFromRequest(ReadRequest.Builder reqBuilder) {
361+
String databaseId = parentChannel.extractDatabaseIdFromSession(reqBuilder.getSession());
362+
ByteString transactionId = transactionIdFromSelector(reqBuilder.getTransaction());
363+
ChannelEndpoint endpoint = parentChannel.affinityEndpoint(transactionId);
364+
ChannelFinder finder = null;
365+
if (databaseId != null) {
366+
finder = parentChannel.getOrCreateChannelFinder(databaseId);
367+
ChannelEndpoint routed = finder.findServer(reqBuilder);
368+
if (endpoint == null) {
369+
endpoint = routed;
370+
}
371+
}
372+
return new RoutingDecision(finder, endpoint);
373+
}
374+
375+
private RoutingDecision routeFromRequest(ExecuteSqlRequest.Builder reqBuilder) {
376+
String databaseId = parentChannel.extractDatabaseIdFromSession(reqBuilder.getSession());
377+
ByteString transactionId = transactionIdFromSelector(reqBuilder.getTransaction());
378+
ChannelEndpoint endpoint = parentChannel.affinityEndpoint(transactionId);
379+
ChannelFinder finder = null;
380+
if (databaseId != null) {
381+
finder = parentChannel.getOrCreateChannelFinder(databaseId);
382+
ChannelEndpoint routed = finder.findServer(reqBuilder);
383+
if (endpoint == null) {
384+
endpoint = routed;
385+
}
386+
}
387+
return new RoutingDecision(finder, endpoint);
388+
}
389+
}
390+
391+
private static final class RoutingDecision {
392+
@Nullable private final ChannelFinder finder;
393+
@Nullable private final ChannelEndpoint endpoint;
394+
395+
private RoutingDecision(@Nullable ChannelFinder finder, @Nullable ChannelEndpoint endpoint) {
396+
this.finder = finder;
397+
this.endpoint = endpoint;
398+
}
376399
}
377400

378401
static final class KeyAwareClientCallListener<ResponseT>

0 commit comments

Comments
 (0)