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

Commit 519f0ef

Browse files
committed
incorporate changes
1 parent 540f387 commit 519f0ef

3 files changed

Lines changed: 207 additions & 23 deletions

File tree

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

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -66,12 +66,7 @@ public void update(CacheUpdate update) {
6666
}
6767

6868
public ChannelEndpoint findServer(ReadRequest.Builder reqBuilder) {
69-
recipeCache.computeKeys(reqBuilder);
70-
return fillRoutingHint(
71-
reqBuilder.getTransaction(),
72-
reqBuilder.getDirectedReadOptions(),
73-
KeyRangeCache.RangeMode.COVERING_SPLIT,
74-
reqBuilder.getRoutingHintBuilder());
69+
return findServer(reqBuilder, preferLeader(reqBuilder.getTransaction()));
7570
}
7671

7772
public ChannelEndpoint findServer(ReadRequest.Builder reqBuilder, boolean preferLeader) {
@@ -84,12 +79,7 @@ public ChannelEndpoint findServer(ReadRequest.Builder reqBuilder, boolean prefer
8479
}
8580

8681
public ChannelEndpoint findServer(ExecuteSqlRequest.Builder reqBuilder) {
87-
recipeCache.computeKeys(reqBuilder);
88-
return fillRoutingHint(
89-
reqBuilder.getTransaction(),
90-
reqBuilder.getDirectedReadOptions(),
91-
KeyRangeCache.RangeMode.PICK_RANDOM,
92-
reqBuilder.getRoutingHintBuilder());
82+
return findServer(reqBuilder, preferLeader(reqBuilder.getTransaction()));
9383
}
9484

9585
public ChannelEndpoint findServer(ExecuteSqlRequest.Builder reqBuilder, boolean preferLeader) {

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

Lines changed: 26 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@
1818

1919
import com.google.api.core.InternalApi;
2020
import com.google.api.gax.grpc.InstantiatingGrpcChannelProvider;
21+
import com.google.common.cache.Cache;
22+
import com.google.common.cache.CacheBuilder;
2123
import com.google.protobuf.ByteString;
2224
import com.google.spanner.v1.BeginTransactionRequest;
2325
import com.google.spanner.v1.CommitRequest;
@@ -51,6 +53,7 @@
5153
*/
5254
@InternalApi
5355
final class KeyAwareChannel extends ManagedChannel {
56+
private static final long MAX_TRACKED_READ_ONLY_TRANSACTIONS = 100_000L;
5457
private static final String STREAMING_READ_METHOD = "google.spanner.v1.Spanner/StreamingRead";
5558
private static final String STREAMING_SQL_METHOD =
5659
"google.spanner.v1.Spanner/ExecuteStreamingSql";
@@ -69,7 +72,9 @@ final class KeyAwareChannel extends ManagedChannel {
6972
private final Map<ByteString, String> transactionAffinities = new ConcurrentHashMap<>();
7073
// Maps read-only transaction IDs to their preferLeader value.
7174
// Strong reads → true (prefer leader), Stale reads → false (any replica).
72-
private final Map<ByteString, Boolean> readOnlyTransactions = new ConcurrentHashMap<>();
75+
// Bounded to prevent unbounded growth if application code does not close read-only transactions.
76+
private final Cache<ByteString, Boolean> readOnlyTxPreferLeader =
77+
CacheBuilder.newBuilder().maximumSize(MAX_TRACKED_READ_ONLY_TRANSACTIONS).build();
7378

7479
private KeyAwareChannel(
7580
InstantiatingGrpcChannelProvider channelProvider,
@@ -187,7 +192,7 @@ private void clearAffinity(ByteString transactionId) {
187192
return;
188193
}
189194
transactionAffinities.remove(transactionId);
190-
readOnlyTransactions.remove(transactionId);
195+
readOnlyTxPreferLeader.invalidate(transactionId);
191196
}
192197

193198
void clearTransactionAffinity(ByteString transactionId) {
@@ -197,22 +202,22 @@ void clearTransactionAffinity(ByteString transactionId) {
197202
private boolean isReadOnlyTransaction(ByteString transactionId) {
198203
return transactionId != null
199204
&& !transactionId.isEmpty()
200-
&& readOnlyTransactions.containsKey(transactionId);
205+
&& readOnlyTxPreferLeader.getIfPresent(transactionId) != null;
201206
}
202207

203208
@Nullable
204209
private Boolean readOnlyPreferLeader(ByteString transactionId) {
205210
if (transactionId == null || transactionId.isEmpty()) {
206211
return null;
207212
}
208-
return readOnlyTransactions.get(transactionId);
213+
return readOnlyTxPreferLeader.getIfPresent(transactionId);
209214
}
210215

211216
private void trackReadOnlyTransaction(ByteString transactionId, boolean preferLeader) {
212217
if (transactionId == null || transactionId.isEmpty()) {
213218
return;
214219
}
215-
readOnlyTransactions.put(transactionId, preferLeader);
220+
readOnlyTxPreferLeader.put(transactionId, preferLeader);
216221
}
217222

218223
private void recordAffinity(
@@ -334,12 +339,14 @@ public void sendMessage(RequestT message) {
334339

335340
if (message instanceof ReadRequest) {
336341
ReadRequest.Builder reqBuilder = ((ReadRequest) message).toBuilder();
342+
maybeTrackReadOnlyBegin(reqBuilder.getTransaction());
337343
RoutingDecision routing = routeFromRequest(reqBuilder);
338344
finder = routing.finder;
339345
endpoint = routing.endpoint;
340346
message = (RequestT) reqBuilder.build();
341347
} else if (message instanceof ExecuteSqlRequest) {
342348
ExecuteSqlRequest.Builder reqBuilder = ((ExecuteSqlRequest) message).toBuilder();
349+
maybeTrackReadOnlyBegin(reqBuilder.getTransaction());
343350
RoutingDecision routing = routeFromRequest(reqBuilder);
344351
finder = routing.finder;
345352
endpoint = routing.endpoint;
@@ -515,6 +522,14 @@ void maybeClearAffinity() {
515522
parentChannel.clearAffinity(transactionIdToClear);
516523
}
517524

525+
private void maybeTrackReadOnlyBegin(TransactionSelector selector) {
526+
if (selector.getSelectorCase() == TransactionSelector.SelectorCase.BEGIN
527+
&& selector.getBegin().hasReadOnly()) {
528+
isReadOnlyBegin = true;
529+
readOnlyIsStrong = selector.getBegin().getReadOnly().getStrong();
530+
}
531+
}
532+
518533
private RoutingDecision routeFromRequest(ReadRequest.Builder reqBuilder) {
519534
String databaseId = parentChannel.extractDatabaseIdFromSession(reqBuilder.getSession());
520535
ByteString transactionId = transactionIdFromSelector(reqBuilder.getTransaction());
@@ -524,14 +539,14 @@ private RoutingDecision routeFromRequest(ReadRequest.Builder reqBuilder) {
524539
ChannelFinder finder = null;
525540
if (databaseId != null) {
526541
finder = parentChannel.getOrCreateChannelFinder(databaseId);
542+
}
543+
if (databaseId != null && endpoint == null) {
527544
Boolean preferLeaderOverride = parentChannel.readOnlyPreferLeader(transactionId);
528545
ChannelEndpoint routed =
529546
preferLeaderOverride != null
530547
? finder.findServer(reqBuilder, preferLeaderOverride)
531548
: finder.findServer(reqBuilder);
532-
if (endpoint == null) {
533-
endpoint = routed;
534-
}
549+
endpoint = routed;
535550
}
536551
return new RoutingDecision(finder, endpoint);
537552
}
@@ -545,14 +560,14 @@ private RoutingDecision routeFromRequest(ExecuteSqlRequest.Builder reqBuilder) {
545560
ChannelFinder finder = null;
546561
if (databaseId != null) {
547562
finder = parentChannel.getOrCreateChannelFinder(databaseId);
563+
}
564+
if (databaseId != null && endpoint == null) {
548565
Boolean preferLeaderOverride = parentChannel.readOnlyPreferLeader(transactionId);
549566
ChannelEndpoint routed =
550567
preferLeaderOverride != null
551568
? finder.findServer(reqBuilder, preferLeaderOverride)
552569
: finder.findServer(reqBuilder);
553-
if (endpoint == null) {
554-
endpoint = routed;
555-
}
570+
endpoint = routed;
556571
}
557572
return new RoutingDecision(finder, endpoint);
558573
}

google-cloud-spanner/src/test/java/com/google/cloud/spanner/spi/v1/KeyAwareChannelTest.java

Lines changed: 179 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -391,6 +391,117 @@ public void readOnlyTransactionRoutesEachReadIndependently() throws Exception {
391391
assertThat(harness.defaultManagedChannel.callCount()).isEqualTo(2);
392392
}
393393

394+
@Test
395+
public void readOnlyInlinedBeginExecuteSqlRoutesSubsequentRequestsIndependently()
396+
throws Exception {
397+
TestHarness harness = createHarness();
398+
ByteString transactionId = ByteString.copyFromUtf8("ro-inline-sql");
399+
400+
seedCache(harness, createTwoRangeCacheUpdate());
401+
402+
// First query begins a read-only transaction inline and routes to server-a.
403+
ClientCall<ExecuteSqlRequest, ResultSet> firstCall =
404+
harness.channel.newCall(SpannerGrpc.getExecuteSqlMethod(), CallOptions.DEFAULT);
405+
firstCall.start(new CapturingListener<ResultSet>(), new Metadata());
406+
firstCall.sendMessage(
407+
ExecuteSqlRequest.newBuilder()
408+
.setSession(SESSION)
409+
.setTransaction(
410+
TransactionSelector.newBuilder()
411+
.setBegin(
412+
TransactionOptions.newBuilder()
413+
.setReadOnly(
414+
TransactionOptions.ReadOnly.newBuilder()
415+
.setReturnReadTimestamp(true)
416+
.build())
417+
.build()))
418+
.setRoutingHint(RoutingHint.newBuilder().setKey(bytes("b")).build())
419+
.build());
420+
421+
assertThat(harness.endpointCache.callCountForAddress("server-a:1234")).isEqualTo(1);
422+
423+
@SuppressWarnings("unchecked")
424+
RecordingClientCall<ExecuteSqlRequest, ResultSet> firstDelegate =
425+
(RecordingClientCall<ExecuteSqlRequest, ResultSet>)
426+
harness.endpointCache.latestCallForAddress("server-a:1234");
427+
firstDelegate.emitOnMessage(
428+
ResultSet.newBuilder()
429+
.setMetadata(
430+
ResultSetMetadata.newBuilder()
431+
.setTransaction(Transaction.newBuilder().setId(transactionId)))
432+
.build());
433+
434+
// Second query in same txn should route by key to server-b, not affinity-pin to server-a.
435+
ClientCall<ExecuteSqlRequest, ResultSet> secondCall =
436+
harness.channel.newCall(SpannerGrpc.getExecuteSqlMethod(), CallOptions.DEFAULT);
437+
secondCall.start(new CapturingListener<ResultSet>(), new Metadata());
438+
secondCall.sendMessage(
439+
ExecuteSqlRequest.newBuilder()
440+
.setSession(SESSION)
441+
.setTransaction(TransactionSelector.newBuilder().setId(transactionId))
442+
.setRoutingHint(RoutingHint.newBuilder().setKey(bytes("n")).build())
443+
.build());
444+
445+
assertThat(harness.endpointCache.callCountForAddress("server-a:1234")).isEqualTo(1);
446+
assertThat(harness.endpointCache.callCountForAddress("server-b:1234")).isEqualTo(1);
447+
assertThat(harness.defaultManagedChannel.callCount()).isEqualTo(1);
448+
}
449+
450+
@Test
451+
public void readOnlyInlinedBeginReadRoutesSubsequentRequestsIndependently() throws Exception {
452+
TestHarness harness = createHarness();
453+
ByteString transactionId = ByteString.copyFromUtf8("ro-inline-read");
454+
455+
seedCache(harness, createTwoRangeCacheUpdate());
456+
457+
// First read begins a read-only transaction inline and routes to server-a.
458+
ClientCall<ReadRequest, PartialResultSet> firstCall =
459+
harness.channel.newCall(SpannerGrpc.getStreamingReadMethod(), CallOptions.DEFAULT);
460+
firstCall.start(new CapturingListener<PartialResultSet>(), new Metadata());
461+
firstCall.sendMessage(
462+
ReadRequest.newBuilder()
463+
.setSession(SESSION)
464+
.setTransaction(
465+
TransactionSelector.newBuilder()
466+
.setBegin(
467+
TransactionOptions.newBuilder()
468+
.setReadOnly(
469+
TransactionOptions.ReadOnly.newBuilder()
470+
.setReturnReadTimestamp(true)
471+
.build())
472+
.build()))
473+
.setRoutingHint(RoutingHint.newBuilder().setKey(bytes("b")).build())
474+
.build());
475+
476+
assertThat(harness.endpointCache.callCountForAddress("server-a:1234")).isEqualTo(1);
477+
478+
@SuppressWarnings("unchecked")
479+
RecordingClientCall<ReadRequest, PartialResultSet> firstDelegate =
480+
(RecordingClientCall<ReadRequest, PartialResultSet>)
481+
harness.endpointCache.latestCallForAddress("server-a:1234");
482+
firstDelegate.emitOnMessage(
483+
PartialResultSet.newBuilder()
484+
.setMetadata(
485+
ResultSetMetadata.newBuilder()
486+
.setTransaction(Transaction.newBuilder().setId(transactionId)))
487+
.build());
488+
489+
// Second read in same txn should route by key to server-b, not affinity-pin to server-a.
490+
ClientCall<ReadRequest, PartialResultSet> secondCall =
491+
harness.channel.newCall(SpannerGrpc.getStreamingReadMethod(), CallOptions.DEFAULT);
492+
secondCall.start(new CapturingListener<PartialResultSet>(), new Metadata());
493+
secondCall.sendMessage(
494+
ReadRequest.newBuilder()
495+
.setSession(SESSION)
496+
.setTransaction(TransactionSelector.newBuilder().setId(transactionId))
497+
.setRoutingHint(RoutingHint.newBuilder().setKey(bytes("n")).build())
498+
.build());
499+
500+
assertThat(harness.endpointCache.callCountForAddress("server-a:1234")).isEqualTo(1);
501+
assertThat(harness.endpointCache.callCountForAddress("server-b:1234")).isEqualTo(1);
502+
assertThat(harness.defaultManagedChannel.callCount()).isEqualTo(1);
503+
}
504+
394505
@Test
395506
public void readOnlyTransactionDoesNotRecordAffinity() throws Exception {
396507
TestHarness harness = createHarness();
@@ -484,6 +595,63 @@ public void readOnlyTransactionCleanupOnClose() throws Exception {
484595
harness.channel.clearTransactionAffinity(transactionId);
485596
}
486597

598+
private static CacheUpdate createTwoRangeCacheUpdate() {
599+
return CacheUpdate.newBuilder()
600+
.setDatabaseId(7L)
601+
.addRange(
602+
Range.newBuilder()
603+
.setStartKey(bytes("a"))
604+
.setLimitKey(bytes("m"))
605+
.setGroupUid(1L)
606+
.setSplitId(1L)
607+
.setGeneration(bytes("1")))
608+
.addRange(
609+
Range.newBuilder()
610+
.setStartKey(bytes("m"))
611+
.setLimitKey(bytes("z"))
612+
.setGroupUid(2L)
613+
.setSplitId(2L)
614+
.setGeneration(bytes("1")))
615+
.addGroup(
616+
Group.newBuilder()
617+
.setGroupUid(1L)
618+
.setGeneration(bytes("1"))
619+
.addTablets(
620+
Tablet.newBuilder()
621+
.setTabletUid(1L)
622+
.setServerAddress("server-a:1234")
623+
.setIncarnation(bytes("1"))
624+
.setDistance(0)))
625+
.addGroup(
626+
Group.newBuilder()
627+
.setGroupUid(2L)
628+
.setGeneration(bytes("1"))
629+
.addTablets(
630+
Tablet.newBuilder()
631+
.setTabletUid(2L)
632+
.setServerAddress("server-b:1234")
633+
.setIncarnation(bytes("1"))
634+
.setDistance(0)))
635+
.build();
636+
}
637+
638+
private static void seedCache(TestHarness harness, CacheUpdate cacheUpdate) {
639+
ClientCall<ExecuteSqlRequest, ResultSet> seedCall =
640+
harness.channel.newCall(SpannerGrpc.getExecuteSqlMethod(), CallOptions.DEFAULT);
641+
seedCall.start(new CapturingListener<ResultSet>(), new Metadata());
642+
seedCall.sendMessage(
643+
ExecuteSqlRequest.newBuilder()
644+
.setSession(SESSION)
645+
.setRoutingHint(RoutingHint.newBuilder().setKey(bytes("a")).build())
646+
.build());
647+
648+
@SuppressWarnings("unchecked")
649+
RecordingClientCall<ExecuteSqlRequest, ResultSet> seedDelegate =
650+
(RecordingClientCall<ExecuteSqlRequest, ResultSet>)
651+
harness.defaultManagedChannel.latestCall();
652+
seedDelegate.emitOnMessage(ResultSet.newBuilder().setCacheUpdate(cacheUpdate).build());
653+
}
654+
487655
private static TestHarness createHarness() throws IOException {
488656
FakeEndpointCache endpointCache = new FakeEndpointCache(DEFAULT_ADDRESS);
489657
InstantiatingGrpcChannelProvider provider =
@@ -574,6 +742,17 @@ int callCountForAddress(String address) {
574742
FakeEndpoint endpoint = endpoints.get(address);
575743
return endpoint == null ? 0 : endpoint.channel.callCount();
576744
}
745+
746+
RecordingClientCall<?, ?> latestCallForAddress(String address) {
747+
if (defaultAddress.equals(address)) {
748+
return defaultEndpoint.channel.latestCall();
749+
}
750+
FakeEndpoint endpoint = endpoints.get(address);
751+
if (endpoint == null) {
752+
throw new IllegalStateException("No endpoint for address: " + address);
753+
}
754+
return endpoint.channel.latestCall();
755+
}
577756
}
578757

579758
private static final class FakeEndpoint implements ChannelEndpoint {

0 commit comments

Comments
 (0)