diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java index a4d7614ebd2b..4df8d54c2385 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java @@ -433,38 +433,47 @@ TransactionSelector getTransactionSelector() { return selector; } - ApiFuture futureToWaitFor = null; - txnLock.lock(); - try { - if (transactionId != null) { - return TransactionSelector.newBuilder().setId(transactionId).build(); + final long deadlineNanos = + System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(WAIT_FOR_INLINE_BEGIN_TIMEOUT_MILLIS); + while (true) { + ApiFuture futureToWaitFor; + txnLock.lock(); + try { + if (transactionId != null) { + return TransactionSelector.newBuilder().setId(transactionId).build(); + } + if (transactionIdFuture == null) { + transactionIdFuture = SettableApiFuture.create(); + return TransactionSelector.newBuilder() + .setBegin(createReadOnlyTransactionOptions()) + .build(); + } + futureToWaitFor = transactionIdFuture; + } finally { + txnLock.unlock(); } - if (transactionIdFuture == null) { - transactionIdFuture = SettableApiFuture.create(); + + try { return TransactionSelector.newBuilder() - .setBegin(createReadOnlyTransactionOptions()) + .setId(futureToWaitFor.get(deadlineNanos - System.nanoTime(), TimeUnit.NANOSECONDS)) .build(); + } catch (ExecutionException e) { + // The statement that carried the inline BeginTransaction failed before a transaction was + // returned, and failTransactionIdFuture has reset transactionIdFuture. Retry the loop so + // this statement either takes over the BeginTransaction itself, or waits for another + // statement that already has. The error is propagated to the failed statement itself and + // should not fail unrelated statements. Read-only transactions take no locks, so unlike + // the read/write inline-begin path there is no need to abort the entire transaction. + } catch (TimeoutException e) { + throw SpannerExceptionFactory.newSpannerException( + ErrorCode.DEADLINE_EXCEEDED, + "Timeout while waiting for an inlined read-only transaction to be returned by another" + + " statement.", + e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw SpannerExceptionFactory.newSpannerExceptionForCancellation(null, e); } - futureToWaitFor = transactionIdFuture; - } finally { - txnLock.unlock(); - } - - try { - return TransactionSelector.newBuilder() - .setId(futureToWaitFor.get(WAIT_FOR_INLINE_BEGIN_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)) - .build(); - } catch (ExecutionException e) { - throw SpannerExceptionFactory.asSpannerException(e.getCause()); - } catch (TimeoutException e) { - throw SpannerExceptionFactory.newSpannerException( - ErrorCode.DEADLINE_EXCEEDED, - "Timeout while waiting for an inlined read-only transaction to be returned by another" - + " statement.", - e); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw SpannerExceptionFactory.newSpannerExceptionForCancellation(null, e); } } @@ -624,6 +633,10 @@ private void failTransactionIdFuture(Throwable t) { try { if (transactionIdFuture != null && !transactionIdFuture.isDone()) { transactionIdFuture.setException(t); + // Reset the future so the next statement attempts a new inline BeginTransaction instead + // of failing on the error of an earlier, unrelated statement. This mirrors the explicit + // BeginTransaction path, where a failed begin is retried by the next statement. + transactionIdFuture = null; } } finally { txnLock.unlock(); diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java index d1779b459e4c..50cb2cc4a842 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java @@ -67,6 +67,7 @@ import java.util.Map; import java.util.TimeZone; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -725,18 +726,32 @@ public void multiUseReadOnlyTransactionCanUseInlineBeginForQuery() throws ParseE } @Test - public void multiUseReadOnlyTransactionInlineBeginFirstQueryErrorPropagates() { + public void multiUseReadOnlyTransactionInlineBeginRetriesBeginAfterFailedFirstStatement() + throws ParseException { SpannerException error = SpannerExceptionFactory.newSpannerException(ErrorCode.INVALID_ARGUMENT, "bad query"); - final ArgumentCaptor consumer = + final ArgumentCaptor queryConsumer = ArgumentCaptor.forClass(SpannerRpc.ResultStreamConsumer.class); - final ArgumentCaptor request = + final ArgumentCaptor queryRequest = ArgumentCaptor.forClass(ExecuteSqlRequest.class); Mockito.when( - rpc.executeQuery(request.capture(), consumer.capture(), anyMap(), any(), eq(false))) + rpc.executeQuery( + queryRequest.capture(), queryConsumer.capture(), anyMap(), any(), eq(false))) .then( invocation -> { - consumer.getValue().onError(error); + queryConsumer.getValue().onError(error); + return new NoOpStreamingCall(); + }); + PartialResultSet readResultSet = inlineBeginResultSet("recovered-tx"); + final ArgumentCaptor readConsumer = + ArgumentCaptor.forClass(SpannerRpc.ResultStreamConsumer.class); + final ArgumentCaptor readRequest = ArgumentCaptor.forClass(ReadRequest.class); + Mockito.when( + rpc.read(readRequest.capture(), readConsumer.capture(), anyMap(), any(), eq(false))) + .then( + invocation -> { + readConsumer.getValue().onPartialResultSet(readResultSet); + readConsumer.getValue().onCompleted(); return new NoOpStreamingCall(); }); @@ -745,18 +760,119 @@ public void multiUseReadOnlyTransactionInlineBeginFirstQueryErrorPropagates() { SpannerException e = assertThrows(SpannerException.class, () -> rs.next()); assertEquals(ErrorCode.INVALID_ARGUMENT, e.getErrorCode()); } + // The transaction should not be poisoned by the failed first statement: the next statement + // should attempt a new inline BeginTransaction. + txn.readRow("Dummy", Key.of(), Collections.singletonList("C")); + txn.readRow("Dummy", Key.of(), Collections.singletonList("C")); + assertEquals( + Timestamp.fromProto(Timestamps.parse("2015-10-01T10:54:20.021Z")), + txn.getReadTimestamp()); + } + + Mockito.verify(rpc, Mockito.never()).beginTransaction(Mockito.any(), anyMap(), eq(false)); + assertEquals(1, queryRequest.getAllValues().size()); + assertThat(queryRequest.getAllValues().get(0).getTransaction().hasBegin()).isTrue(); + assertEquals(2, readRequest.getAllValues().size()); + assertThat(readRequest.getAllValues().get(0).getTransaction().hasBegin()).isTrue(); + assertEquals( + ByteString.copyFromUtf8("recovered-tx"), + readRequest.getAllValues().get(1).getTransaction().getId()); + } + + @Test + public void multiUseReadOnlyTransactionInlineBeginRecoversWhenNoTransactionIsReturned() + throws ParseException { + PartialResultSet resultSetWithoutTransaction = resultSetWithoutTransaction(); + PartialResultSet resultSetWithTransaction = inlineBeginResultSet("second-attempt-tx"); + final ArgumentCaptor request = ArgumentCaptor.forClass(ReadRequest.class); + final AtomicInteger callCount = new AtomicInteger(); + Mockito.when(rpc.read(request.capture(), Mockito.any(), anyMap(), any(), eq(false))) + .then( + invocation -> { + SpannerRpc.ResultStreamConsumer consumer = invocation.getArgument(1); + consumer.onPartialResultSet( + callCount.incrementAndGet() == 1 + ? resultSetWithoutTransaction + : resultSetWithTransaction); + consumer.onCompleted(); + return new NoOpStreamingCall(); + }); + + try (ReadOnlyTransaction txn = inlineReadOnlyTransaction()) { SpannerException e = assertThrows( SpannerException.class, () -> txn.readRow("Dummy", Key.of(), Collections.singletonList("C"))); - assertEquals(ErrorCode.INVALID_ARGUMENT, e.getErrorCode()); + assertEquals(ErrorCode.FAILED_PRECONDITION, e.getErrorCode()); + // The next statement should attempt a new inline BeginTransaction instead of failing on the + // error of the first statement. + txn.readRow("Dummy", Key.of(), Collections.singletonList("C")); } Mockito.verify(rpc, Mockito.never()).beginTransaction(Mockito.any(), anyMap(), eq(false)); - assertEquals(1, request.getAllValues().size()); + assertEquals(2, request.getAllValues().size()); assertThat(request.getAllValues().get(0).getTransaction().hasBegin()).isTrue(); - Mockito.verify(rpc, Mockito.never()) - .read(Mockito.any(), Mockito.any(), anyMap(), any(), eq(false)); + assertThat(request.getAllValues().get(1).getTransaction().hasBegin()).isTrue(); + } + + @Test + public void multiUseReadOnlyTransactionInlineBeginConcurrentStatementTakesOverFailedBegin() + throws Exception { + PartialResultSet takeOverResultSet = inlineBeginResultSet("take-over-tx"); + SpannerException error = + SpannerExceptionFactory.newSpannerException(ErrorCode.INVALID_ARGUMENT, "bad read"); + final List requests = Collections.synchronizedList(new ArrayList<>()); + final List consumers = + Collections.synchronizedList(new ArrayList<>()); + final AtomicInteger callCount = new AtomicInteger(); + final CountDownLatch firstRpcStarted = new CountDownLatch(1); + Mockito.when(rpc.read(Mockito.any(), Mockito.any(), anyMap(), any(), eq(false))) + .then( + invocation -> { + int call = callCount.incrementAndGet(); + ReadRequest readRequest = invocation.getArgument(0); + SpannerRpc.ResultStreamConsumer consumer = invocation.getArgument(1); + requests.add(readRequest); + consumers.add(consumer); + if (call == 1) { + firstRpcStarted.countDown(); + } else { + consumer.onPartialResultSet(takeOverResultSet); + consumer.onCompleted(); + } + return new NoOpStreamingCall(); + }); + + ExecutorService executor = Executors.newFixedThreadPool(2); + try (ReadOnlyTransaction txn = inlineReadOnlyTransaction()) { + Future first = + executor.submit(() -> txn.readRow("Dummy", Key.of(), Collections.singletonList("C"))); + assertThat(firstRpcStarted.await(5, TimeUnit.SECONDS)).isTrue(); + + Future second = + executor.submit(() -> txn.readRow("Dummy", Key.of(), Collections.singletonList("C"))); + Thread.sleep(100L); + assertThat(callCount.get()).isEqualTo(1); + assertThat(second.isDone()).isFalse(); + + // Fail the statement that carried the inline BeginTransaction. The concurrently waiting + // statement should take over the BeginTransaction instead of failing with the error of the + // first statement. + consumers.get(0).onError(error); + + ExecutionException e = + assertThrows(ExecutionException.class, () -> first.get(5, TimeUnit.SECONDS)); + assertThat(e.getCause()).isInstanceOf(SpannerException.class); + assertEquals(ErrorCode.INVALID_ARGUMENT, ((SpannerException) e.getCause()).getErrorCode()); + assertThat(second.get(5, TimeUnit.SECONDS)).isNull(); + } finally { + executor.shutdownNow(); + } + + Mockito.verify(rpc, Mockito.never()).beginTransaction(Mockito.any(), anyMap(), eq(false)); + assertEquals(2, requests.size()); + assertThat(requests.get(0).getTransaction().hasBegin()).isTrue(); + assertThat(requests.get(1).getTransaction().hasBegin()).isTrue(); } @Test