From 12ce287a5ca2be808c793a431251c28bbba8ea12 Mon Sep 17 00:00:00 2001 From: Lawrence Qiu Date: Mon, 6 Jul 2026 18:32:48 +0000 Subject: [PATCH 1/3] fix(bigquerystorage): fix flaky ConnectionWorkerPoolTest Fix ConnectionWorkerPoolTest.testMultiStreamAppend_appendWhileClosing by limiting max connections per region to 5. This prevents the pool from scaling up to 6 connections due to timing-dependent overwhelm of mock connections, when the test expects exactly 5. Also fix a concurrency issue in FakeBigQueryWriteImpl where requestReceivedInstants ArrayList was accessed concurrently without synchronization, causing ArrayIndexOutOfBoundsException (reported as UNKNOWN RPC error in tests). Fixes https://github.com/googleapis/google-cloud-java/issues/13637 TAG=agy CONV=79e14d4e-c762-4dad-a0d6-de0dd0e3d3f6 --- .../bigquery/storage/v1/ConnectionWorkerPoolTest.java | 4 +++- .../cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java | 8 ++++++-- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerPoolTest.java b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerPoolTest.java index 01de76298955..1fc07927cb7b 100644 --- a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerPoolTest.java +++ b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerPoolTest.java @@ -273,8 +273,10 @@ void testMultiStreamClosed_multiplexingEnabled() throws Exception { @Test void testMultiStreamAppend_appendWhileClosing() throws Exception { + // Limit max connections to 5 (same as min connections) to prevent timing-dependent + // scale up during concurrent appends in this test, which expects exactly 5 connections. ConnectionWorkerPool.setOptions( - Settings.builder().setMaxConnectionsPerRegion(10).setMinConnectionsPerRegion(5).build()); + Settings.builder().setMaxConnectionsPerRegion(5).setMinConnectionsPerRegion(5).build()); ConnectionWorkerPool connectionWorkerPool = createConnectionWorkerPool( /* maxRequests= */ 3, /* maxBytes= */ 100000, java.time.Duration.ofSeconds(5)); diff --git a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java index d8cbd758b0e1..4205e80708d4 100644 --- a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java +++ b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java @@ -114,7 +114,9 @@ public String toString() { } public ArrayList getLatestRequestReceivedInstants() { - return requestReceivedInstants; + synchronized (requestReceivedInstants) { + return new ArrayList<>(requestReceivedInstants); + } } @Override @@ -203,7 +205,9 @@ public StreamObserver appendRows( new StreamObserver() { @Override public void onNext(AppendRowsRequest value) { - requestReceivedInstants.add(Instant.now()); + synchronized (requestReceivedInstants) { + requestReceivedInstants.add(Instant.now()); + } recordCount++; requests.add(value); long offset = value.getOffset().getValue(); From ce1714717781343c5cd33f3ce9f0120670f45e93 Mon Sep 17 00:00:00 2001 From: Lawrence Qiu Date: Mon, 6 Jul 2026 18:44:49 +0000 Subject: [PATCH 2/3] refactor(bigquerystorage): make FakeBigQueryWriteImpl thread-safe Improve thread safety of FakeBigQueryWriteImpl by: - Replacing requestReceivedInstants ArrayList with CopyOnWriteArrayList and removing manual synchronization. - Replacing recordCount, connectionCount, expectedOffset (long) and responseIndex (int) with AtomicLong/AtomicInteger to avoid race conditions and lost updates. - Synchronizing the offset verification block to ensure atomic check-and-increment of expectedOffset. This addresses feedback on PR #13662 regarding potential concurrency issues in the test fake. TAG=agy CONV=79e14d4e-c762-4dad-a0d6-de0dd0e3d3f6 --- .../storage/v1/FakeBigQueryWriteImpl.java | 80 ++++++++++--------- 1 file changed, 43 insertions(+), 37 deletions(-) diff --git a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java index 4205e80708d4..291ac80d3626 100644 --- a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java +++ b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java @@ -27,11 +27,13 @@ import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.Supplier; import java.util.logging.Logger; @@ -59,11 +61,11 @@ class FakeBigQueryWriteImpl extends BigQueryWriteGrpc.BigQueryWriteImplBase { private long numberTimesToClose = 0; private long closeAfter = 0; - private long recordCount = 0; - private long connectionCount = 0; + private final AtomicLong recordCount = new AtomicLong(0); + private final AtomicLong connectionCount = new AtomicLong(0); private long closeForeverAfter = 0; - private int responseIndex = 0; - private long expectedOffset = 0; + private final AtomicInteger responseIndex = new AtomicInteger(0); + private final AtomicLong expectedOffset = new AtomicLong(0); private boolean verifyOffset = false; private boolean returnErrorDuringExclusiveStreamRetry = false; private boolean returnErrorUntilRetrySuccess = false; @@ -74,7 +76,8 @@ class FakeBigQueryWriteImpl extends BigQueryWriteGrpc.BigQueryWriteImplBase { private final Map, Boolean> connectionToFirstRequest = new ConcurrentHashMap<>(); private Status failedStatus = Status.ABORTED; - private ArrayList requestReceivedInstants = new ArrayList<>(); + private final CopyOnWriteArrayList requestReceivedInstants = + new CopyOnWriteArrayList<>(); /** Class used to save the state of a possible response. */ public static class Response { @@ -114,9 +117,7 @@ public String toString() { } public ArrayList getLatestRequestReceivedInstants() { - synchronized (requestReceivedInstants) { - return new ArrayList<>(requestReceivedInstants); - } + return new ArrayList<>(requestReceivedInstants); } @Override @@ -155,7 +156,7 @@ void waitForResponseScheduled() throws InterruptedException { /* Return the number of times the stream was connected. */ public long getConnectionCount() { - return connectionCount; + return connectionCount.get(); } void setFailedStatus(Status failedStatus) { @@ -199,22 +200,23 @@ private Response determineResponse(long offset) { @Override public StreamObserver appendRows( final StreamObserver responseObserver) { - this.connectionCount++; + connectionCount.incrementAndGet(); connectionToFirstRequest.put(responseObserver, true); StreamObserver requestObserver = new StreamObserver() { @Override public void onNext(AppendRowsRequest value) { - synchronized (requestReceivedInstants) { - requestReceivedInstants.add(Instant.now()); - } - recordCount++; + requestReceivedInstants.add(Instant.now()); + long currentRecordCount = recordCount.incrementAndGet(); + int currentResponseIndex = responseIndex.getAndIncrement(); + int responseIndexAfterIncrement = currentResponseIndex + 1; + long currentConnectionCount = connectionCount.get(); + requests.add(value); long offset = value.getOffset().getValue(); if (offset == -1 || !value.hasOffset()) { - offset = responseIndex; + offset = currentResponseIndex; } - responseIndex++; if (responseSleep.compareTo(Duration.ZERO) > 0) { LOG.info("Sleeping before response for " + responseSleep.toString()); Uninterruptibles.sleepUninterruptibly( @@ -237,33 +239,37 @@ public void onNext(AppendRowsRequest value) { } connectionToFirstRequest.put(responseObserver, false); if (closeAfter > 0 - && responseIndex % closeAfter == 0 - && recordCount % closeAfter == 0 - && (numberTimesToClose == 0 || connectionCount <= numberTimesToClose)) { + && responseIndexAfterIncrement % closeAfter == 0 + && currentRecordCount % closeAfter == 0 + && (numberTimesToClose == 0 || currentConnectionCount <= numberTimesToClose)) { LOG.info("Shutting down connection from test..."); responseObserver.onError(failedStatus.asException()); - } else if (closeForeverAfter > 0 && recordCount > closeForeverAfter) { + } else if (closeForeverAfter > 0 && currentRecordCount > closeForeverAfter) { LOG.info("Shutting down connection from test..."); responseObserver.onError(failedStatus.asException()); } else { Response response = determineResponse(offset); - if (verifyOffset - && !response.getResponse().hasError() - && response.getResponse().getAppendResult().getOffset().getValue() > -1) { - // No error and offset is present; verify order - if (response.getResponse().getAppendResult().getOffset().getValue() - != expectedOffset) { - com.google.rpc.Status status = - com.google.rpc.Status.newBuilder().setCode(Code.INTERNAL_VALUE).build(); - response = new Response(AppendRowsResponse.newBuilder().setError(status).build()); - } else { - LOG.info( - String.format( - "asserted offset: %s expected: %s", - response.getResponse().getAppendResult().getOffset().getValue(), - expectedOffset)); - LOG.info(String.format("sending response: %s", response.getResponse())); - expectedOffset++; + if (verifyOffset) { + synchronized (expectedOffset) { + if (!response.getResponse().hasError() + && response.getResponse().getAppendResult().getOffset().getValue() > -1) { + long currentExpectedOffset = expectedOffset.get(); + if (response.getResponse().getAppendResult().getOffset().getValue() + != currentExpectedOffset) { + com.google.rpc.Status status = + com.google.rpc.Status.newBuilder().setCode(Code.INTERNAL_VALUE).build(); + response = + new Response(AppendRowsResponse.newBuilder().setError(status).build()); + } else { + LOG.info( + String.format( + "asserted offset: %s expected: %s", + response.getResponse().getAppendResult().getOffset().getValue(), + currentExpectedOffset)); + LOG.info(String.format("sending response: %s", response.getResponse())); + expectedOffset.incrementAndGet(); + } + } } } sendResponse(response, responseObserver); From 32527a9cca1dd3d8abed83f30c85a7efd10ea3c3 Mon Sep 17 00:00:00 2001 From: Lawrence Qiu Date: Mon, 6 Jul 2026 18:48:24 +0000 Subject: [PATCH 3/3] refactor(bigquerystorage): use lock-free compareAndSet for expectedOffset verification Replace the synchronized block in FakeBigQueryWriteImpl offset verification with AtomicLong.compareAndSet. This makes the check-and-increment atomic and lock-free, improving performance and thread safety for concurrent streams. Added comments explaining the approach. TAG=agy CONV=79e14d4e-c762-4dad-a0d6-de0dd0e3d3f6 --- .../storage/v1/FakeBigQueryWriteImpl.java | 43 ++++++++++--------- 1 file changed, 22 insertions(+), 21 deletions(-) diff --git a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java index 291ac80d3626..ce39e7fbe744 100644 --- a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java +++ b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/FakeBigQueryWriteImpl.java @@ -249,27 +249,28 @@ public void onNext(AppendRowsRequest value) { responseObserver.onError(failedStatus.asException()); } else { Response response = determineResponse(offset); - if (verifyOffset) { - synchronized (expectedOffset) { - if (!response.getResponse().hasError() - && response.getResponse().getAppendResult().getOffset().getValue() > -1) { - long currentExpectedOffset = expectedOffset.get(); - if (response.getResponse().getAppendResult().getOffset().getValue() - != currentExpectedOffset) { - com.google.rpc.Status status = - com.google.rpc.Status.newBuilder().setCode(Code.INTERNAL_VALUE).build(); - response = - new Response(AppendRowsResponse.newBuilder().setError(status).build()); - } else { - LOG.info( - String.format( - "asserted offset: %s expected: %s", - response.getResponse().getAppendResult().getOffset().getValue(), - currentExpectedOffset)); - LOG.info(String.format("sending response: %s", response.getResponse())); - expectedOffset.incrementAndGet(); - } - } + if (verifyOffset + && !response.getResponse().hasError() + && response.getResponse().getAppendResult().getOffset().getValue() > -1) { + long responseOffset = + response.getResponse().getAppendResult().getOffset().getValue(); + // Atomically verify that the response offset matches the expected offset + // and increment the expected offset for the next request. This avoids + // using a synchronized block while ensuring thread safety across concurrent + // streams. + if (!expectedOffset.compareAndSet(responseOffset, responseOffset + 1)) { + LOG.info( + String.format( + "Offset mismatch: expected %s, got %s", + expectedOffset.get(), responseOffset)); + com.google.rpc.Status status = + com.google.rpc.Status.newBuilder().setCode(Code.INTERNAL_VALUE).build(); + response = new Response(AppendRowsResponse.newBuilder().setError(status).build()); + } else { + LOG.info( + String.format( + "asserted offset: %s expected: %s", responseOffset, responseOffset)); + LOG.info(String.format("sending response: %s", response.getResponse())); } } sendResponse(response, responseObserver);