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 ae18a1e4f5f..27b6440ab9b 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 @@ -451,11 +451,6 @@ void initTransaction() { this.tracer = builder.tracer; } - @Override - public void setSpan(ISpan span) { - this.span = span; - } - long getSeqNo() { return seqNo.incrementAndGet(); } diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncTransactionManagerImpl.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncTransactionManagerImpl.java index 8b20dd824a0..cd83b247518 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncTransactionManagerImpl.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AsyncTransactionManagerImpl.java @@ -47,11 +47,6 @@ final class AsyncTransactionManagerImpl this.options = Options.fromTransactionOptions(options); } - @Override - public void setSpan(ISpan span) { - this.span = span; - } - @Override public void close() { closeAsync(); diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseClientImpl.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseClientImpl.java index e9c9818f451..0dab9639c74 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseClientImpl.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/DatabaseClientImpl.java @@ -21,6 +21,7 @@ import com.google.cloud.spanner.Options.TransactionOption; import com.google.cloud.spanner.Options.UpdateOption; import com.google.cloud.spanner.SessionPool.PooledSessionFuture; +import com.google.cloud.spanner.SessionPool.SessionFuture; import com.google.cloud.spanner.SpannerImpl.ClosedException; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Function; @@ -52,6 +53,10 @@ PooledSessionFuture getSession() { return pool.getSession(); } + SessionFuture getMultiplexedSession() { + return pool.getMultiplexedSessionWithFallback(); + } + @Override public Dialect getDialect() { return pool.getDialect(); @@ -123,7 +128,7 @@ public ServerStream batchWriteAtLeastOnce( public ReadContext singleUse() { ISpan span = tracer.spanBuilder(READ_ONLY_TRANSACTION); try (IScope s = tracer.withSpan(span)) { - return getSession().singleUse(); + return getMultiplexedSession().singleUse(); } catch (RuntimeException e) { span.setStatus(e); span.end(); @@ -135,7 +140,7 @@ public ReadContext singleUse() { public ReadContext singleUse(TimestampBound bound) { ISpan span = tracer.spanBuilder(READ_ONLY_TRANSACTION); try (IScope s = tracer.withSpan(span)) { - return getSession().singleUse(bound); + return getMultiplexedSession().singleUse(bound); } catch (RuntimeException e) { span.setStatus(e); span.end(); @@ -147,7 +152,7 @@ public ReadContext singleUse(TimestampBound bound) { public ReadOnlyTransaction singleUseReadOnlyTransaction() { ISpan span = tracer.spanBuilder(READ_ONLY_TRANSACTION); try (IScope s = tracer.withSpan(span)) { - return getSession().singleUseReadOnlyTransaction(); + return getMultiplexedSession().singleUseReadOnlyTransaction(); } catch (RuntimeException e) { span.setStatus(e); span.end(); @@ -159,7 +164,7 @@ public ReadOnlyTransaction singleUseReadOnlyTransaction() { public ReadOnlyTransaction singleUseReadOnlyTransaction(TimestampBound bound) { ISpan span = tracer.spanBuilder(READ_ONLY_TRANSACTION); try (IScope s = tracer.withSpan(span)) { - return getSession().singleUseReadOnlyTransaction(bound); + return getMultiplexedSession().singleUseReadOnlyTransaction(bound); } catch (RuntimeException e) { span.setStatus(e); span.end(); @@ -171,7 +176,7 @@ public ReadOnlyTransaction singleUseReadOnlyTransaction(TimestampBound bound) { public ReadOnlyTransaction readOnlyTransaction() { ISpan span = tracer.spanBuilder(READ_ONLY_TRANSACTION); try (IScope s = tracer.withSpan(span)) { - return getSession().readOnlyTransaction(); + return getMultiplexedSession().readOnlyTransaction(); } catch (RuntimeException e) { span.setStatus(e); span.end(); @@ -183,7 +188,7 @@ public ReadOnlyTransaction readOnlyTransaction() { public ReadOnlyTransaction readOnlyTransaction(TimestampBound bound) { ISpan span = tracer.spanBuilder(READ_ONLY_TRANSACTION); try (IScope s = tracer.withSpan(span)) { - return getSession().readOnlyTransaction(bound); + return getMultiplexedSession().readOnlyTransaction(bound); } catch (RuntimeException e) { span.setStatus(e); span.end(); diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/PartitionedDmlTransaction.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/PartitionedDmlTransaction.java index d498bb232a1..2b1be0ad9cd 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/PartitionedDmlTransaction.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/PartitionedDmlTransaction.java @@ -138,10 +138,6 @@ public void invalidate() { isValid = false; } - /** No-op method needed to implement SessionTransaction interface. */ - @Override - public void setSpan(ISpan span) {} - private Duration tryUpdateTimeout(final Duration timeout, final Stopwatch stopwatch) { final Duration remainingTimeout = timeout.minus(stopwatch.elapsed(TimeUnit.MILLISECONDS), ChronoUnit.MILLIS); 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 8c4a0068599..be2e7ff30b2 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 @@ -91,9 +91,6 @@ interface SessionTransaction { /** Invalidates the transaction, generally because a new one has been started on the session. */ void invalidate(); - - /** Registers the current span on the transaction. */ - void setSpan(ISpan span); } private final SpannerImpl spanner; @@ -105,7 +102,6 @@ interface SessionTransaction { private volatile Instant lastUseTime; @Nullable private final Instant createTime; private final boolean isMultiplexed; - private ISpan currentSpan; SessionImpl(SpannerImpl spanner, String name, Map options) { this.spanner = spanner; @@ -143,14 +139,6 @@ public String getName() { return options; } - void setCurrentSpan(ISpan span) { - currentSpan = span; - } - - ISpan getCurrentSpan() { - return currentSpan; - } - Instant getLastUseTime() { return lastUseTime; } @@ -308,7 +296,7 @@ public ReadContext singleUse(TimestampBound bound) { .setDefaultPrefetchChunks(spanner.getDefaultPrefetchChunks()) .setDefaultDecodeMode(spanner.getDefaultDecodeMode()) .setDefaultDirectedReadOptions(spanner.getOptions().getDirectedReadOptions()) - .setSpan(currentSpan) + .setSpan(tracer.getCurrentSpan()) .setTracer(tracer) .setExecutorProvider(spanner.getAsyncExecutorProvider()) .build()); @@ -330,7 +318,7 @@ public ReadOnlyTransaction singleUseReadOnlyTransaction(TimestampBound bound) { .setDefaultPrefetchChunks(spanner.getDefaultPrefetchChunks()) .setDefaultDecodeMode(spanner.getDefaultDecodeMode()) .setDefaultDirectedReadOptions(spanner.getOptions().getDirectedReadOptions()) - .setSpan(currentSpan) + .setSpan(tracer.getCurrentSpan()) .setTracer(tracer) .setExecutorProvider(spanner.getAsyncExecutorProvider()) .buildSingleUseReadOnlyTransaction()); @@ -352,7 +340,7 @@ public ReadOnlyTransaction readOnlyTransaction(TimestampBound bound) { .setDefaultPrefetchChunks(spanner.getDefaultPrefetchChunks()) .setDefaultDecodeMode(spanner.getDefaultDecodeMode()) .setDefaultDirectedReadOptions(spanner.getOptions().getDirectedReadOptions()) - .setSpan(currentSpan) + .setSpan(tracer.getCurrentSpan()) .setTracer(tracer) .setExecutorProvider(spanner.getAsyncExecutorProvider()) .build()); @@ -370,12 +358,12 @@ public AsyncRunner runAsync(TransactionOption... options) { @Override public TransactionManager transactionManager(TransactionOption... options) { - return new TransactionManagerImpl(this, currentSpan, tracer, options); + return new TransactionManagerImpl(this, tracer.getCurrentSpan(), tracer, options); } @Override public AsyncTransactionManagerImpl transactionManagerAsync(TransactionOption... options) { - return new AsyncTransactionManagerImpl(this, currentSpan, options); + return new AsyncTransactionManagerImpl(this, tracer.getCurrentSpan(), options); } @Override @@ -470,7 +458,7 @@ TransactionContextImpl newTransaction(Options options) { .setDefaultQueryOptions(spanner.getDefaultQueryOptions(databaseId)) .setDefaultPrefetchChunks(spanner.getDefaultPrefetchChunks()) .setDefaultDecodeMode(spanner.getDefaultDecodeMode()) - .setSpan(currentSpan) + .setSpan(tracer.getCurrentSpan()) .setTracer(tracer) .setExecutorProvider(spanner.getAsyncExecutorProvider()) .setClock(poolMaintainerClock == null ? new Clock() : poolMaintainerClock) @@ -479,14 +467,12 @@ TransactionContextImpl newTransaction(Options options) { T setActive(@Nullable T ctx) { throwIfTransactionsPending(); - - if (activeTransaction != null) { - activeTransaction.invalidate(); - } - activeTransaction = ctx; - readyTransactionId = null; - if (activeTransaction != null) { - activeTransaction.setSpan(currentSpan); + if (!this.isMultiplexed) { + if (activeTransaction != null) { + activeTransaction.invalidate(); + } + activeTransaction = ctx; + readyTransactionId = null; } return ctx; } 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 a5eee9d0db7..ba0d17b2690 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 @@ -1908,15 +1908,11 @@ public void prepareReadWriteTransaction() { private void keepAlive() { markUsed(); - final ISpan previousSpan = delegate.getCurrentSpan(); - delegate.setCurrentSpan(tracer.getBlankSpan()); try (ResultSet resultSet = delegate .singleUse(TimestampBound.ofMaxStaleness(60, TimeUnit.SECONDS)) .executeQuery(Statement.newBuilder("SELECT 1").build())) { resultSet.next(); - } finally { - delegate.setCurrentSpan(previousSpan); } } @@ -1955,7 +1951,6 @@ public SessionImpl getDelegate() { @Override public void markBusy(ISpan span) { - this.delegate.setCurrentSpan(span); this.state = SessionState.BUSY; } diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionManagerImpl.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionManagerImpl.java index 95ffd1168b2..9e245148e61 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionManagerImpl.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionManagerImpl.java @@ -40,15 +40,6 @@ final class TransactionManagerImpl implements TransactionManager, SessionTransac this.options = Options.fromTransactionOptions(options); } - ISpan getSpan() { - return span; - } - - @Override - public void setSpan(ISpan span) { - this.span = span; - } - @Override public TransactionContext begin() { Preconditions.checkState(txn == null, "begin can only be called once"); diff --git a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java index 370d2f662f5..e1021b81046 100644 --- a/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java +++ b/google-cloud-spanner/src/main/java/com/google/cloud/spanner/TransactionRunnerImpl.java @@ -985,11 +985,7 @@ public TransactionRunner allowNestedTransaction() { this.options = Options.fromTransactionOptions(options); this.txn = session.newTransaction(this.options); this.tracer = session.getTracer(); - } - - @Override - public void setSpan(ISpan span) { - this.span = span; + this.span = tracer.getCurrentSpan(); } @Nullable diff --git a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java index ff50bf52e70..fbd5c5d020e 100644 --- a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java +++ b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java @@ -140,7 +140,6 @@ public void setUp() { Span oTspan = mock(Span.class); ISpan span = new OpenTelemetrySpan(oTspan); when(oTspan.makeCurrent()).thenReturn(mock(Scope.class)); - ((SessionImpl) session).setCurrentSpan(span); // We expect the same options, "options", on all calls on "session". options = optionsCaptor.getValue(); } 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 e1c295c14de..da6fae3faab 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 @@ -1455,7 +1455,6 @@ public void testSessionNotFoundReadWriteTransaction() { when(closedSession.beginTransactionAsync(any(), eq(true))).thenThrow(sessionNotFound); when(closedSession.getTracer()).thenReturn(tracer); TransactionRunnerImpl closedTransactionRunner = new TransactionRunnerImpl(closedSession); - closedTransactionRunner.setSpan(span); when(closedSession.readWriteTransaction()).thenReturn(closedTransactionRunner); final SessionImpl openSession = mock(SessionImpl.class); @@ -1470,7 +1469,6 @@ public void testSessionNotFoundReadWriteTransaction() { .thenReturn(ApiFutures.immediateFuture(ByteString.copyFromUtf8("open-txn"))); when(openSession.getTracer()).thenReturn(tracer); TransactionRunnerImpl openTransactionRunner = new TransactionRunnerImpl(openSession); - openTransactionRunner.setSpan(span); when(openSession.readWriteTransaction()).thenReturn(openTransactionRunner); ResultSet openResultSet = mock(ResultSet.class); diff --git a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java index cf5f901f65e..ff0a6921450 100644 --- a/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java +++ b/google-cloud-spanner/src/test/java/com/google/cloud/spanner/TransactionRunnerImplTest.java @@ -146,7 +146,6 @@ public void setUp() { Span oTspan = mock(Span.class); span = new OpenTelemetrySpan(oTspan); when(oTspan.makeCurrent()).thenReturn(mock(Scope.class)); - transactionRunner.setSpan(span); } @SuppressWarnings("unchecked") @@ -312,9 +311,7 @@ public void prepareReadWriteTransaction() { throw new IllegalStateException(); } }; - session.setCurrentSpan(new OpenTelemetrySpan(mock(io.opentelemetry.api.trace.Span.class))); TransactionRunnerImpl runner = new TransactionRunnerImpl(session); - runner.setSpan(span); assertThat(usedInlinedBegin).isFalse(); runner.run( transaction -> { @@ -347,7 +344,6 @@ private long[] batchDmlException(int status) { ApiFutures.immediateFuture(ByteString.copyFromUtf8(UUID.randomUUID().toString()))); when(session.getName()).thenReturn(SessionId.of("p", "i", "d", "test").getName()); TransactionRunnerImpl runner = new TransactionRunnerImpl(session); - runner.setSpan(span); ExecuteBatchDmlResponse response1 = ExecuteBatchDmlResponse.newBuilder() .addResultSets(