diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java index f2ea7af64e5..4d17ba4e1b7 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java @@ -809,6 +809,7 @@ public final void invalidate() { @Override public void close() { + session.onTransactionDone(); span.end(); synchronized (lock) { isClosed = true; diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DelayedReadContext.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DelayedReadContext.java index 976f611fe70..752c41cdea3 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DelayedReadContext.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DelayedReadContext.java @@ -130,7 +130,14 @@ public ResultSet analyzeQuery(Statement statement, QueryAnalyzeMode queryMode) { } @Override - public void close() {} + public void close() { + try { + this.readContextFuture.get().close(); + } catch (Throwable ignore) { + // Ignore any errors during close, as this error has already propagated to the user through + // other means. + } + } /** * Represents a {@link ReadContext} using a multiplexed session that is not yet ready. The diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java index ae7301afae9..e742481be2c 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClient.java @@ -39,6 +39,7 @@ import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; /** @@ -63,6 +64,8 @@ static class MultiplexedSessionTransaction extends SessionImpl { private final int singleUseChannelHint; + private boolean done; + MultiplexedSessionTransaction( MultiplexedSessionDatabaseClient client, ISpan span, @@ -73,6 +76,7 @@ static class MultiplexedSessionTransaction extends SessionImpl { this.client = client; this.singleUse = singleUse; this.singleUseChannelHint = singleUseChannelHint; + this.client.numSessionsAcquired.incrementAndGet(); setCurrentSpan(span); } @@ -103,6 +107,20 @@ void onReadDone() { } } + @Override + void onTransactionDone() { + boolean markedDone = false; + synchronized (this) { + if (!this.done) { + this.done = true; + markedDone = true; + } + } + if (markedDone) { + client.numSessionsReleased.incrementAndGet(); + } + } + @Override public void close() { // no-op, we don't want to delete the multiplexed session. @@ -152,6 +170,10 @@ public void close() { private final AtomicReference resourceNotFoundException = new AtomicReference<>(); + private final AtomicLong numSessionsAcquired = new AtomicLong(); + + private final AtomicLong numSessionsReleased = new AtomicLong(); + /** * This flag is set to true if the server return UNIMPLEMENTED when we try to create a multiplexed * session. TODO: Remove once this is guaranteed to be available. @@ -239,6 +261,14 @@ boolean isValid() { return resourceNotFoundException.get() == null; } + AtomicLong getNumSessionsAcquired() { + return this.numSessionsAcquired; + } + + AtomicLong getNumSessionsReleased() { + return this.numSessionsReleased; + } + boolean isMultiplexedSessionsSupported() { return !this.unimplemented.get(); } diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionImpl.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionImpl.java index 8e2c545c8b9..ab985cebf45 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionImpl.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionImpl.java @@ -484,6 +484,8 @@ void onError(SpannerException spannerException) {} void onReadDone() {} + void onTransactionDone() {} + TraceWrapper getTracer() { return tracer; } diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionPool.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionPool.java index fa58e7799d4..3277e04f963 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionPool.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SessionPool.java @@ -68,7 +68,6 @@ import com.google.common.base.Function; import com.google.common.base.MoreObjects; import com.google.common.base.Preconditions; -import com.google.common.base.Supplier; import com.google.common.collect.ImmutableList; import com.google.common.util.concurrent.ForwardingListenableFuture; import com.google.common.util.concurrent.ForwardingListenableFuture.SimpleForwardingListenableFuture; @@ -156,7 +155,8 @@ void maybeWaitOnMinSessions() { } } - private abstract static class CachedResultSetSupplier implements Supplier { + private abstract static class CachedResultSetSupplier + implements com.google.common.base.Supplier { private ResultSet cached; @@ -2265,7 +2265,6 @@ public String getName() { @Override public void close() { synchronized (lock) { - numMultiplexedSessionsReleased++; if (lastException != null && isDatabaseOrInstanceNotFound(lastException)) { SessionPool.this.resourceNotFoundException = MoreObjects.firstNonNull( @@ -2771,15 +2770,9 @@ enum Position { @GuardedBy("lock") private long numSessionsAcquired = 0; - @GuardedBy("lock") - private long numMultiplexedSessionsAcquired = 0; - @GuardedBy("lock") private long numSessionsReleased = 0; - @GuardedBy("lock") - private long numMultiplexedSessionsReleased = 0; - @GuardedBy("lock") private long numIdleSessionsRemoved = 0; @@ -2830,7 +2823,9 @@ static SessionPool createPool( SessionClient sessionClient, TraceWrapper tracer, List labelValues, - Attributes attributes) { + Attributes attributes, + AtomicLong numMultiplexedSessionsAcquired, + AtomicLong numMultiplexedSessionsReleased) { final SessionPoolOptions sessionPoolOptions = spannerOptions.getSessionPoolOptions(); // A clock instance is passed in {@code SessionPoolOptions} in order to allow mocking via tests. @@ -2846,7 +2841,9 @@ static SessionPool createPool( tracer, labelValues, spannerOptions.getOpenTelemetry(), - attributes); + attributes, + numMultiplexedSessionsAcquired, + numMultiplexedSessionsReleased); } static SessionPool createPool( @@ -2884,7 +2881,9 @@ static SessionPool createPool( tracer, SPANNER_DEFAULT_LABEL_VALUES, openTelemetry, - null); + null, + new AtomicLong(), + new AtomicLong()); } static SessionPool createPool( @@ -2898,7 +2897,9 @@ static SessionPool createPool( TraceWrapper tracer, List labelValues, OpenTelemetry openTelemetry, - Attributes attributes) { + Attributes attributes, + AtomicLong numMultiplexedSessionsAcquired, + AtomicLong numMultiplexedSessionsReleased) { SessionPool pool = new SessionPool( poolOptions, @@ -2912,7 +2913,9 @@ static SessionPool createPool( tracer, labelValues, openTelemetry, - attributes); + attributes, + numMultiplexedSessionsAcquired, + numMultiplexedSessionsReleased); pool.initPool(); return pool; } @@ -2929,7 +2932,9 @@ private SessionPool( TraceWrapper tracer, List labelValues, OpenTelemetry openTelemetry, - Attributes attributes) { + Attributes attributes, + AtomicLong numMultiplexedSessionsAcquired, + AtomicLong numMultiplexedSessionsReleased) { this.options = options; this.databaseRole = databaseRole; this.executorFactory = executorFactory; @@ -2940,8 +2945,13 @@ private SessionPool( this.initialReleasePosition = initialReleasePosition; this.poolMaintainer = new PoolMaintainer(); this.tracer = tracer; - this.initOpenCensusMetricsCollection(metricRegistry, labelValues); - this.initOpenTelemetryMetricsCollection(openTelemetry, attributes); + this.initOpenCensusMetricsCollection( + metricRegistry, + labelValues, + numMultiplexedSessionsAcquired, + numMultiplexedSessionsReleased); + this.initOpenTelemetryMetricsCollection( + openTelemetry, attributes, numMultiplexedSessionsAcquired, numMultiplexedSessionsReleased); this.waitOnMinSessionsLatch = options.getMinSessions() > 0 ? new CountDownLatch(1) : new CountDownLatch(0); this.waitOnMultiplexedSessionsLatch = new CountDownLatch(1); @@ -3143,7 +3153,7 @@ boolean isValid() { /** * Returns a multiplexed session. The method fallbacks to a regular session if {@link - * SessionPoolOptions#useMultiplexedSession} is not set. + * SessionPoolOptions#getUseMultiplexedSession} is not set. */ SessionFutureWrapper getMultiplexedSessionWithFallback() throws SpannerException { if (useMultiplexedSessions()) { @@ -3250,8 +3260,6 @@ private void incrementNumSessionsInUse(boolean isMultiplexed) { maxSessionsInUse = numSessionsInUse; } numSessionsAcquired++; - } else { - numMultiplexedSessionsAcquired++; } } } @@ -3775,7 +3783,10 @@ public void onSessionCreateFailure(Throwable t, int createFailureForSessionCount * exporter, it allows users to monitor client behavior. */ private void initOpenCensusMetricsCollection( - MetricRegistry metricRegistry, List labelValues) { + MetricRegistry metricRegistry, + List labelValues, + AtomicLong numMultiplexedSessionsAcquired, + AtomicLong numMultiplexedSessionsReleased) { if (!SpannerOptions.isEnabledOpenCensusMetrics()) { return; } @@ -3860,18 +3871,14 @@ private void initOpenCensusMetricsCollection( labelValuesWithRegularSessions, this, sessionPool -> sessionPool.numSessionsAcquired); numAcquiredSessionsMetric.removeTimeSeries(labelValuesWithMultiplexedSessions); numAcquiredSessionsMetric.createTimeSeries( - labelValuesWithMultiplexedSessions, - this, - sessionPool -> sessionPool.numMultiplexedSessionsAcquired); + labelValuesWithMultiplexedSessions, this, unused -> numMultiplexedSessionsAcquired.get()); numReleasedSessionsMetric.removeTimeSeries(labelValuesWithRegularSessions); numReleasedSessionsMetric.createTimeSeries( labelValuesWithRegularSessions, this, sessionPool -> sessionPool.numSessionsReleased); numReleasedSessionsMetric.removeTimeSeries(labelValuesWithMultiplexedSessions); numReleasedSessionsMetric.createTimeSeries( - labelValuesWithMultiplexedSessions, - this, - sessionPool -> sessionPool.numMultiplexedSessionsReleased); + labelValuesWithMultiplexedSessions, this, unused -> numMultiplexedSessionsReleased.get()); List labelValuesWithBeingPreparedType = new ArrayList<>(labelValues); labelValuesWithBeingPreparedType.add(NUM_SESSIONS_BEING_PREPARED); @@ -3909,7 +3916,10 @@ private void initOpenCensusMetricsCollection( * an exporter, it allows users to monitor client behavior. */ private void initOpenTelemetryMetricsCollection( - OpenTelemetry openTelemetry, Attributes attributes) { + OpenTelemetry openTelemetry, + Attributes attributes, + AtomicLong numMultiplexedSessionsAcquired, + AtomicLong numMultiplexedSessionsReleased) { if (openTelemetry == null || !SpannerOptions.isEnabledOpenTelemetryMetrics()) { return; } @@ -3981,7 +3991,8 @@ private void initOpenTelemetryMetricsCollection( .buildWithCallback( measurement -> { measurement.record(this.numSessionsAcquired, attributesRegularSession); - measurement.record(this.numMultiplexedSessionsAcquired, attributesMultiplexedSession); + measurement.record( + numMultiplexedSessionsAcquired.get(), attributesMultiplexedSession); }); meter @@ -3991,7 +4002,8 @@ private void initOpenTelemetryMetricsCollection( .buildWithCallback( measurement -> { measurement.record(this.numSessionsReleased, attributesRegularSession); - measurement.record(this.numMultiplexedSessionsReleased, attributesMultiplexedSession); + measurement.record( + numMultiplexedSessionsReleased.get(), attributesMultiplexedSession); }); } } diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerImpl.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerImpl.java index 3c782d269cf..86b5de01c69 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerImpl.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/SpannerImpl.java @@ -50,6 +50,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; import java.util.logging.Level; import java.util.logging.Logger; import javax.annotation.Nullable; @@ -271,17 +272,29 @@ public DatabaseClient getDatabaseClient(DatabaseId db) { attributesBuilder.put("database", db.getDatabase()); attributesBuilder.put("instance_id", db.getInstanceId().getName()); + boolean useMultiplexedSession = + getOptions().getSessionPoolOptions().getUseMultiplexedSession(); + MultiplexedSessionDatabaseClient multiplexedSessionDatabaseClient = + useMultiplexedSession + ? new MultiplexedSessionDatabaseClient(SpannerImpl.this.getSessionClient(db)) + : null; + AtomicLong numMultiplexedSessionsAcquired = + useMultiplexedSession + ? multiplexedSessionDatabaseClient.getNumSessionsAcquired() + : new AtomicLong(); + AtomicLong numMultiplexedSessionsReleased = + useMultiplexedSession + ? multiplexedSessionDatabaseClient.getNumSessionsReleased() + : new AtomicLong(); SessionPool pool = SessionPool.createPool( getOptions(), SpannerImpl.this.getSessionClient(db), this.tracer, labelValues, - attributesBuilder.build()); - MultiplexedSessionDatabaseClient multiplexedSessionDatabaseClient = - getOptions().getSessionPoolOptions().getUseMultiplexedSession() - ? new MultiplexedSessionDatabaseClient(SpannerImpl.this.getSessionClient(db)) - : null; + attributesBuilder.build(), + numMultiplexedSessionsAcquired, + numMultiplexedSessionsReleased); pool.maybeWaitOnMinSessions(); DatabaseClientImpl dbClient = createDatabaseClient(clientId, pool, multiplexedSessionDatabaseClient); diff --git a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/DatabaseClientImplTest.java b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/DatabaseClientImplTest.java index b0b75a2b996..2c29bbe80e0 100644 --- a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/DatabaseClientImplTest.java +++ b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/DatabaseClientImplTest.java @@ -3859,12 +3859,12 @@ public void testCreateSessionsFailure_shouldNotPropagateToCloseMethod() { // Simulate session creation failures on the backend. mockSpanner.setCreateSessionExecutionTime( SimulatedExecutionTime.ofStickyException(Status.RESOURCE_EXHAUSTED.asRuntimeException())); - DatabaseClient client = - spannerWithEmptySessionPool.getDatabaseClient( - DatabaseId.of(TEST_PROJECT, TEST_INSTANCE, TEST_DATABASE)); // This will not cause any failure as getting a session from the pool is guaranteed to be // non-blocking, and any exceptions will be delayed until actual query execution. mockSpanner.freeze(); + DatabaseClient client = + spannerWithEmptySessionPool.getDatabaseClient( + DatabaseId.of(TEST_PROJECT, TEST_INSTANCE, TEST_DATABASE)); try (ResultSet rs = client.singleUse().executeQuery(SELECT1)) { mockSpanner.unfreeze(); SpannerException e = assertThrows(SpannerException.class, rs::next); diff --git a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java index e8cfe0d61ce..bf4a02a10c5 100644 --- a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java +++ b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/MultiplexedSessionDatabaseClientMockServerTest.java @@ -100,6 +100,10 @@ public void testMultiUseReadOnlyTransactionUsesSameSession() { List requests = mockSpanner.getRequestsOfType(ExecuteSqlRequest.class); assertEquals(2, requests.size()); assertEquals(requests.get(0).getSession(), requests.get(1).getSession()); + + assertNotNull(client.multiplexedSessionDatabaseClient); + assertEquals(1L, client.multiplexedSessionDatabaseClient.getNumSessionsAcquired().get()); + assertEquals(1L, client.multiplexedSessionDatabaseClient.getNumSessionsReleased().get()); } @Test @@ -129,6 +133,10 @@ public void testNewTransactionUsesNewSession() { List requests = mockSpanner.getRequestsOfType(ExecuteSqlRequest.class); assertEquals(2, requests.size()); assertNotEquals(requests.get(0).getSession(), requests.get(1).getSession()); + + assertNotNull(client.multiplexedSessionDatabaseClient); + assertEquals(2L, client.multiplexedSessionDatabaseClient.getNumSessionsAcquired().get()); + assertEquals(2L, client.multiplexedSessionDatabaseClient.getNumSessionsReleased().get()); } @Test @@ -165,6 +173,12 @@ public void testMaintainerMaintainsMultipleClients() { Set sessionIds = requests.stream().map(ExecuteSqlRequest::getSession).collect(Collectors.toSet()); assertEquals(4, sessionIds.size()); + + for (DatabaseClientImpl client : ImmutableList.of(client1, client2)) { + assertNotNull(client.multiplexedSessionDatabaseClient); + assertEquals(2L, client.multiplexedSessionDatabaseClient.getNumSessionsAcquired().get()); + assertEquals(2L, client.multiplexedSessionDatabaseClient.getNumSessionsReleased().get()); + } } @Test @@ -196,6 +210,10 @@ public void testUnimplementedErrorOnCreation_fallsBackToRegularSessions() { Session session = mockSpanner.getSession(requests.get(0).getSession()); assertNotNull(session); assertFalse(session.getMultiplexed()); + + assertNotNull(client.multiplexedSessionDatabaseClient); + assertEquals(0L, client.multiplexedSessionDatabaseClient.getNumSessionsAcquired().get()); + assertEquals(0L, client.multiplexedSessionDatabaseClient.getNumSessionsReleased().get()); } @Test @@ -234,6 +252,10 @@ public void testUnimplementedErrorOnCreation_fallsBackToRegularSessions() { Session session = mockSpanner.getSession(requests.get(0).getSession()); assertNotNull(session); assertFalse(session.getMultiplexed()); + + assertNotNull(client.multiplexedSessionDatabaseClient); + assertEquals(0L, client.multiplexedSessionDatabaseClient.getNumSessionsAcquired().get()); + assertEquals(0L, client.multiplexedSessionDatabaseClient.getNumSessionsReleased().get()); } @Test @@ -281,6 +303,10 @@ public void testMaintainerInvalidatesMultiplexedSessionClientIfUnimplemented() { Session session2 = mockSpanner.getSession(requests.get(1).getSession()); assertNotNull(session2); assertFalse(session2.getMultiplexed()); + + assertNotNull(client.multiplexedSessionDatabaseClient); + assertEquals(1L, client.multiplexedSessionDatabaseClient.getNumSessionsAcquired().get()); + assertEquals(1L, client.multiplexedSessionDatabaseClient.getNumSessionsReleased().get()); } private void waitForSessionToBeReplaced(DatabaseClientImpl client) { diff --git a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionPoolTest.java b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionPoolTest.java index 75580c25124..ab7eb80cf90 100644 --- a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionPoolTest.java +++ b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionPoolTest.java @@ -117,6 +117,7 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.logging.Level; import java.util.logging.Logger; import java.util.stream.Collectors; @@ -192,7 +193,9 @@ private SessionPool createPool( tracer, labelValues, OpenTelemetry.noop(), - null); + null, + new AtomicLong(), + new AtomicLong()); } private SessionPool createPool( @@ -212,7 +215,9 @@ private SessionPool createPool( tracer, labelValues, openTelemetry, - attributes); + attributes, + new AtomicLong(), + new AtomicLong()); } @BeforeClass diff --git a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AllTypesMockServerTest.java b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AllTypesMockServerTest.java index 0ecc837f4f5..3313fa53426 100644 --- a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AllTypesMockServerTest.java +++ b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/connection/AllTypesMockServerTest.java @@ -27,7 +27,6 @@ import com.google.cloud.spanner.SingerProto.Genre; import com.google.cloud.spanner.SingerProto.SingerInfo; import com.google.cloud.spanner.Statement; -import com.google.common.base.Stopwatch; import com.google.common.collect.ImmutableList; import com.google.protobuf.ListValue; import com.google.protobuf.NullValue; @@ -41,9 +40,6 @@ import java.util.Base64; import java.util.List; import java.util.Map; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import java.util.stream.IntStream; import org.junit.After; @@ -555,28 +551,6 @@ public void clearRequests() { mockSpanner.clearRequests(); } - @Test - public void testCounter() throws InterruptedException { - ExecutorService executor = Executors.newFixedThreadPool(1); - Stopwatch watch = Stopwatch.createStarted(); - for (int n = 0; n < 256; n++) { - executor.submit( - () -> { - try (Connection connection = createConnection()) { - connection.setAutocommit(true); - for (int i = 0; i < 100; i++) { - try (ResultSet resultSet = connection.executeQuery(SELECT_STATEMENT)) { - while (resultSet.next()) {} - } - } - } - }); - } - executor.shutdown(); - assertTrue(executor.awaitTermination(60L, TimeUnit.SECONDS)); - System.out.println("Elapsed: " + watch.elapsed()); - } - @Test public void testSelectAllTypes() { try (Connection connection = createConnection()) {