Skip to content
This repository was archived by the owner on Aug 30, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
f579642
feat(spanner): support multiplexed session for Partitioned read or qu…
pratickchokhani Aug 1, 2024
8c3bd2f
chore(spanner): lint fixes
pratickchokhani Aug 1, 2024
b208044
feat(spanner): support multiplexed session for Partitioned DML operat…
pratickchokhani Aug 2, 2024
e36a6bb
Merge branch 'main' of https://github.com/googleapis/java-spanner int…
pratickchokhani Nov 22, 2024
938ef53
lint(spanner): javadoc fixes.
pratickchokhani Nov 25, 2024
df8b4c0
feat(spanner): Updated unit tests of Partitioned operations for Multi…
pratickchokhani Nov 25, 2024
af625f4
feat(spanner): Updated unit tests of Partitioned operations for Multi…
pratickchokhani Nov 25, 2024
9b99037
Merge branch 'partitioned-query' of github.com:pratickchokhani/java-s…
pratickchokhani Nov 25, 2024
8bde6f8
Merge branch 'partitioned-query' of github.com:pratickchokhani/java-s…
pratickchokhani Nov 25, 2024
44f1105
Merge branch 'partitioned-query' of github.com:pratickchokhani/java-s…
pratickchokhani Nov 25, 2024
00a2578
lint(spanner): Apply suggestions from code review
pratickchokhani Nov 27, 2024
13b94e3
lint(spanner): Apply suggestions from code review
pratickchokhani Nov 27, 2024
8c85310
Merge branch 'partitioned-query' of github.com:pratickchokhani/java-s…
pratickchokhani Nov 27, 2024
36fd5d5
Merge branch 'partitioned-query' of github.com:pratickchokhani/java-s…
pratickchokhani Nov 27, 2024
61de125
Merge branch 'partitioned-query' of github.com:pratickchokhani/java-s…
pratickchokhani Dec 5, 2024
3a66903
Merge branch 'main' of https://github.com/googleapis/java-spanner int…
pratickchokhani Dec 5, 2024
f4271cc
feat(spanner): Modified BatchClientImpl to store multiplexed session …
pratickchokhani Dec 9, 2024
2db1cc6
Merge branch 'main' of https://github.com/googleapis/java-spanner int…
pratickchokhani Dec 9, 2024
263901f
feat(spanner): Removed env variable for Partitioned Ops ensuring that…
pratickchokhani Dec 9, 2024
def751b
lint(spanner): Removed unused variables.
pratickchokhani Dec 9, 2024
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
Merge branch 'main' of https://github.com/googleapis/java-spanner int…
…o partitioned-query
  • Loading branch information
pratickchokhani committed Nov 22, 2024
commit e36a6bb727e8408cf4a8076ffbd9405867ca6f4e
Original file line number Diff line number Diff line change
Expand Up @@ -50,24 +50,4 @@ public ServerStream<BatchWriteResponse> batchWriteAtLeastOnce(
throws SpannerException {
throw new UnsupportedOperationException();
}

@Override
public TransactionRunner readWriteTransaction(TransactionOption... options) {
throw new UnsupportedOperationException();
}

@Override
public TransactionManager transactionManager(TransactionOption... options) {
throw new UnsupportedOperationException();
}

@Override
public AsyncRunner runAsync(TransactionOption... options) {
throw new UnsupportedOperationException();
}

@Override
public AsyncTransactionManager transactionManagerAsync(TransactionOption... options) {
throw new UnsupportedOperationException();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,8 @@ class DatabaseClientImpl implements DatabaseClient {
private final TraceWrapper tracer;
@VisibleForTesting final String clientId;
@VisibleForTesting final SessionPool pool;
final boolean useMultiplexedSessionPartitionedOps;
@VisibleForTesting final MultiplexedSessionDatabaseClient multiplexedSessionDatabaseClient;
final boolean useMultiplexedSessionPartitionedOps;
@VisibleForTesting final boolean useMultiplexedSessionForRW;

final boolean useMultiplexedSessionBlindWrite;
Expand All @@ -46,19 +46,23 @@ class DatabaseClientImpl implements DatabaseClient {
this(
"",
pool,
/* multiplexedSessionDatabaseClient= */ null,
/* useMultiplexedSessionPartitionedOps= */ false,
tracer);
/* useMultiplexedSessionBlindWrite = */ false,
/* multiplexedSessionDatabaseClient = */ null,
/* useMultiplexedSessionPartitionedOps= */ false,
tracer,
/* useMultiplexedSessionForRW = */ false);
}

@VisibleForTesting
DatabaseClientImpl(String clientId, SessionPool pool, TraceWrapper tracer) {
this(
clientId,
pool,
/* multiplexedSessionDatabaseClient= */ null,
/* useMultiplexedSessionPartitionedOps= */ false,
tracer);
/* useMultiplexedSessionBlindWrite = */ false,
/* multiplexedSessionDatabaseClient = */ null,
/* useMultiplexedSessionPartitionedOps= */ false,
tracer,
/* useMultiplexedSessionForRW = */ false);
}

DatabaseClientImpl(
Expand All @@ -67,7 +71,8 @@ class DatabaseClientImpl implements DatabaseClient {
boolean useMultiplexedSessionBlindWrite,
@Nullable MultiplexedSessionDatabaseClient multiplexedSessionDatabaseClient,
boolean useMultiplexedSessionPartitionedOps,
TraceWrapper tracer) {
TraceWrapper tracer,
boolean useMultiplexedSessionForRW) {
this.clientId = clientId;
this.pool = pool;
this.useMultiplexedSessionBlindWrite = useMultiplexedSessionBlindWrite;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import com.google.cloud.spanner.DelayedReadContext.DelayedReadOnlyTransaction;
import com.google.cloud.spanner.MultiplexedSessionDatabaseClient.MultiplexedSessionTransaction;
import com.google.cloud.spanner.Options.UpdateOption;
import com.google.cloud.spanner.Options.TransactionOption;
import com.google.common.util.concurrent.MoreExecutors;
import java.util.concurrent.ExecutionException;

Expand Down Expand Up @@ -123,6 +124,108 @@ public ReadOnlyTransaction readOnlyTransaction(TimestampBound bound) {
MoreExecutors.directExecutor()));
}

/**
* This is a blocking method, as the interface that it implements is also defined as a blocking
* method.
*/
@Override
public CommitResponse writeAtLeastOnceWithOptions(
Iterable<Mutation> mutations, TransactionOption... options) throws SpannerException {
SessionReference sessionReference = getSessionReference();
try (MultiplexedSessionTransaction transaction =
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ true)) {
return transaction.writeAtLeastOnceWithOptions(mutations, options);
}
}

// This is a blocking method, as the interface that it implements is also defined as a blocking
// method.
@Override
public Timestamp write(Iterable<Mutation> mutations) throws SpannerException {
SessionReference sessionReference = getSessionReference();
try (MultiplexedSessionTransaction transaction =
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ false)) {
return transaction.write(mutations);
}
}

// This is a blocking method, as the interface that it implements is also defined as a blocking
// method.
@Override
public CommitResponse writeWithOptions(Iterable<Mutation> mutations, TransactionOption... options)
throws SpannerException {
SessionReference sessionReference = getSessionReference();
try (MultiplexedSessionTransaction transaction =
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ false)) {
return transaction.writeWithOptions(mutations, options);
}
}

@Override
public TransactionRunner readWriteTransaction(TransactionOption... options) {
return new DelayedTransactionRunner(
ApiFutures.transform(
this.sessionFuture,
sessionReference ->
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ false)
.readWriteTransaction(options),
MoreExecutors.directExecutor()));
}

@Override
public TransactionManager transactionManager(TransactionOption... options) {
return new DelayedTransactionManager(
ApiFutures.transform(
this.sessionFuture,
sessionReference ->
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ false)
.transactionManager(options),
MoreExecutors.directExecutor()));
}

@Override
public AsyncRunner runAsync(TransactionOption... options) {
return new DelayedAsyncRunner(
ApiFutures.transform(
this.sessionFuture,
sessionReference ->
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ false)
.runAsync(options),
MoreExecutors.directExecutor()));
}

@Override
public AsyncTransactionManager transactionManagerAsync(TransactionOption... options) {
return new DelayedAsyncTransactionManager(
ApiFutures.transform(
this.sessionFuture,
sessionReference ->
new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ false)
.transactionManagerAsync(options),
MoreExecutors.directExecutor()));
}

/**
* Gets the session reference that this delayed transaction is waiting for. This method should
* only be called by methods that are allowed to be blocking.
*/
private SessionReference getSessionReference() {
try {
return this.sessionFuture.get();
} catch (ExecutionException executionException) {
throw SpannerExceptionFactory.causeAsRunTimeException(executionException);
} catch (InterruptedException interruptedException) {
throw SpannerExceptionFactory.propagateInterrupt(interruptedException);
}
}

/**
* Execute `stmt` within PARTITIONED_DML transaction using multiplexed session. This method is a
* blocking call as the interface expects to return the output of the `stmt`.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
import com.google.api.core.ApiFuture;
import com.google.api.core.ApiFutures;
import com.google.api.core.SettableApiFuture;
import com.google.cloud.Timestamp;
import com.google.cloud.spanner.Options.TransactionOption;
import com.google.cloud.spanner.Options.UpdateOption;
import com.google.cloud.spanner.SessionClient.SessionConsumer;
import com.google.cloud.spanner.SpannerException.ResourceNotFoundException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,14 @@ public class SessionPoolOptions {

private final boolean useMultiplexedSession;

/**
* Controls whether multiplexed session is enabled for blind write or not. This is only used for
* systest soak. TODO: Remove when multiplexed session for blind write is released.
*/
private final boolean useMultiplexedSessionBlindWrite;

private final boolean useMultiplexedSessionForRW;

private final boolean useMultiplexedSessionForPartitionedOps;

// TODO: Change to use java.time.Duration.
Expand Down Expand Up @@ -110,6 +118,14 @@ private SessionPoolOptions(Builder builder) {
(useMultiplexedSessionFromEnvVariable != null)
? useMultiplexedSessionFromEnvVariable
: builder.useMultiplexedSession;
this.useMultiplexedSessionBlindWrite = builder.useMultiplexedSessionBlindWrite;
// useMultiplexedSessionForRW priority => Environment var > private setter > client default
Boolean useMultiplexedSessionForRWFromEnvVariable =
getUseMultiplexedSessionForRWFromEnvVariable();
this.useMultiplexedSessionForRW =
(useMultiplexedSessionForRWFromEnvVariable != null)
? useMultiplexedSessionForRWFromEnvVariable
: builder.useMultiplexedSessionForRW;
// useMultiplexedSessionPartitionedOps priority => Environment var > private setter > client
// default
Boolean useMultiplexedSessionFromEnvVariablePartitionedOps =
Expand Down Expand Up @@ -320,6 +336,20 @@ public boolean getUseMultiplexedSession() {
return useMultiplexedSession;
}

@VisibleForTesting
@InternalApi
protected boolean getUseMultiplexedSessionBlindWrite() {
return getUseMultiplexedSession() && useMultiplexedSessionBlindWrite;
}

@VisibleForTesting
@InternalApi
public boolean getUseMultiplexedSessionForRW() {
// Multiplexed sessions for R/W are enabled only if both global multiplexed sessions and
// read-write multiplexed session flags are set to true.
return getUseMultiplexedSession() && useMultiplexedSessionForRW;
}

@VisibleForTesting
@InternalApi
public boolean getUseMultiplexedSessionPartitionedOps() {
Expand Down Expand Up @@ -558,6 +588,15 @@ public static class Builder {
// Set useMultiplexedSession to true to make multiplexed session the default.
private boolean useMultiplexedSession = false;

// TODO: Remove when multiplexed session for blind write is released.
private boolean useMultiplexedSessionBlindWrite = false;

// This field controls the default behavior of session management for RW operations in Java
// client.
// Set useMultiplexedSessionForRW to true to make multiplexed session for RW operations the
// default.
private boolean useMultiplexedSessionForRW = false;

// This field controls the default behavior of session management in Java client.
// Set useMultiplexedSessionPartitionedOps to true to make multiplexed session the default.
private boolean useMultiplexedSessionPartitionedOps = false;
Expand Down Expand Up @@ -603,6 +642,8 @@ private Builder(SessionPoolOptions options) {
this.randomizePositionQPSThreshold = options.randomizePositionQPSThreshold;
this.inactiveTransactionRemovalOptions = options.inactiveTransactionRemovalOptions;
this.useMultiplexedSession = options.useMultiplexedSession;
this.useMultiplexedSessionBlindWrite = options.useMultiplexedSessionBlindWrite;
this.useMultiplexedSessionForRW = options.useMultiplexedSessionForRW;
this.useMultiplexedSessionPartitionedOps = options.useMultiplexedSessionForPartitionedOps;
this.multiplexedSessionMaintenanceDuration = options.multiplexedSessionMaintenanceDuration;
this.poolMaintainerClock = options.poolMaintainerClock;
Expand Down Expand Up @@ -791,6 +832,28 @@ Builder setUseMultiplexedSession(boolean useMultiplexedSession) {
return this;
}

/**
* This method enables multiplexed sessions for blind writes. This method will be removed in the
* future when multiplexed sessions has been made the default for all operations.
*/
@InternalApi
@VisibleForTesting
Builder setUseMultiplexedSessionBlindWrite(boolean useMultiplexedSessionBlindWrite) {
this.useMultiplexedSessionBlindWrite = useMultiplexedSessionBlindWrite;
return this;
}

/**
* Sets whether the client should use multiplexed session for R/W operations or not. This method
* is intentionally package-private and intended for internal use.
*/
@InternalApi
@VisibleForTesting
Builder setUseMultiplexedSessionForRW(boolean useMultiplexedSessionForRW) {
this.useMultiplexedSessionForRW = useMultiplexedSessionForRW;
return this;
}

/**
* Sets whether the client should use multiplexed session or not. If set to true, the client
* optimises and runs multiple applicable requests concurrently on a single session. A single
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -306,8 +306,10 @@ public DatabaseClient getDatabaseClient(DatabaseId db) {
createDatabaseClient(
clientId,
pool,
getOptions().getSessionPoolOptions().getUseMultiplexedSessionBlindWrite(),
multiplexedSessionDatabaseClient,
getOptions().getSessionPoolOptions().getUseMultiplexedSessionPartitionedOps());
getOptions().getSessionPoolOptions().getUseMultiplexedSessionPartitionedOps(),
useMultiplexedSessionForRW);
dbClients.put(db, dbClient);
return dbClient;
}
Expand All @@ -318,10 +320,18 @@ public DatabaseClient getDatabaseClient(DatabaseId db) {
DatabaseClientImpl createDatabaseClient(
String clientId,
SessionPool pool,
boolean useMultiplexedSessionBlindWrite,
@Nullable MultiplexedSessionDatabaseClient multiplexedSessionClient,
boolean useMultiplexedSessionPartitionedOps) {
boolean useMultiplexedSessionPartitionedOps,
boolean useMultiplexedSessionForRW) {
return new DatabaseClientImpl(
clientId, pool, multiplexedSessionClient, useMultiplexedSessionPartitionedOps, tracer);
clientId,
pool,
useMultiplexedSessionBlindWrite,
multiplexedSessionClient,
useMultiplexedSessionPartitionedOps,
tracer,
useMultiplexedSessionForRW);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,10 @@ private static class SpannerWithClosedSessionsImpl extends SpannerImpl {
DatabaseClientImpl createDatabaseClient(
String clientId,
SessionPool pool,
boolean useMultiplexedSessionBlindWriteIgnore,
MultiplexedSessionDatabaseClient ignore,
boolean useMultiplexedSessionPartitionedOpsIgnore) {
boolean useMultiplexedSessionPartitionedOpsIgnore,
boolean useMultiplexedSessionForRWIgnore) {
return new DatabaseClientWithClosedSessionImpl(clientId, pool, tracer);
}
}
Expand Down
You are viewing a condensed version of this merge commit. You can view the full changes here.