diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncResultSetImpl.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncResultSetImpl.java index 3dd5724532b3..36523d7eb027 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncResultSetImpl.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncResultSetImpl.java @@ -372,9 +372,7 @@ public void run() { try { while (!stop && hasNext) { try { - synchronized (monitor) { - stop = state.shouldStop; - } + stop = shouldStopProducing(); if (!stop) { while (buffer.remainingCapacity() == 0 && !stop) { waitIfPaused(); @@ -389,9 +387,7 @@ public void run() { Math.min(buffer.size() / 2 + 1, buffer.size()), MAX_WAIT_FOR_BUFFER_CONSUMPTION)); bufferConsumptionLatch.await(); - synchronized (monitor) { - stop = state.shouldStop; - } + stop = shouldStopProducing(); } } if (!stop) { @@ -450,6 +446,16 @@ public void run() { } } + private boolean shouldStopProducing() { + synchronized (monitor) { + // A callback that throws leaves the state at CONSUMING, unlike DONE and CANCELLED, so + // shouldStop alone would not catch it. The callback runner is never dispatched again and + // bufferConsumptionLatch is left at zero, so the producer would otherwise spin on a full + // buffer that nothing will drain. + return state.shouldStop || cursorReturnedDoneOrException; + } + } + private void waitIfPaused() throws InterruptedException { CountDownLatch pause; synchronized (monitor) { diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/AsyncResultSetImplTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/AsyncResultSetImplTest.java index 23180356cbfa..3b9c20e044eb 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/AsyncResultSetImplTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/AsyncResultSetImplTest.java @@ -494,6 +494,37 @@ public void callbackReturnsError() throws InterruptedException { } } + @Test + public void callbackThrowsExceptionWhileBufferIsFull_shouldStopProducer() throws Exception { + ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + ResultSet delegate = mock(ResultSet.class); + // Return more rows than the buffer can hold, so the producer is waiting for the callback to + // consume rows from the full buffer when the callback throws an exception. + when(delegate.next()).thenReturn(true); + when(delegate.getCurrentRowAsStruct()).thenReturn(mock(Struct.class)); + final AtomicInteger callbackCounter = new AtomicInteger(); + try (AsyncResultSetImpl rs = new AsyncResultSetImpl(simpleProvider, delegate, 1)) { + rs.setCallback( + executor, + resultSet -> { + callbackCounter.incrementAndGet(); + throw new RuntimeException("async test"); + }); + ExecutionException e = + assertThrows(ExecutionException.class, () -> rs.getResult().get(10L, TimeUnit.SECONDS)); + assertThat(e.getCause()).isInstanceOf(SpannerException.class); + SpannerException se = (SpannerException) e.getCause(); + assertThat(se.getErrorCode()).isEqualTo(ErrorCode.UNKNOWN); + assertThat(se.getMessage()).contains("async test"); + assertThat(callbackCounter.get()).isEqualTo(1); + Mockito.verify(delegate).close(); + } + } finally { + executor.shutdown(); + } + } + @Test public void callbackReturnsDoneBeforeEnd_shouldStopIteration() throws Exception { Executor executor = Executors.newSingleThreadExecutor();