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 2 commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -79,9 +79,7 @@ public BatchReadOnlyTransaction batchReadOnlyTransaction(TimestampBound bound) {
@Override
public BatchReadOnlyTransaction batchReadOnlyTransaction(BatchTransactionId batchTransactionId) {
SessionImpl session =
sessionClient.sessionWithId(
checkNotNull(batchTransactionId).getSessionId(),
batchTransactionId.isMultiplexedSession());
sessionClient.sessionWithId(checkNotNull(batchTransactionId).getSessionId());
return new BatchReadOnlyTransactionImpl(
MultiUseReadOnlyTransaction.newBuilder()
.setSession(session)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@ public class BatchTransactionId implements Serializable {
private final ByteString transactionId;
private final String sessionId;
private final Timestamp timestamp;
private final boolean isMultiplexedSession;
private static final long serialVersionUID = 8067099123096783939L;

BatchTransactionId(
Expand All @@ -43,7 +42,6 @@ public class BatchTransactionId implements Serializable {
this.transactionId = Preconditions.checkNotNull(transactionId);
this.sessionId = Preconditions.checkNotNull(sessionId);
this.timestamp = Preconditions.checkNotNull(timestamp);
this.isMultiplexedSession = isMultiplexedSession;
}

ByteString getTransactionId() {
Expand All @@ -58,15 +56,11 @@ Timestamp getTimestamp() {
return timestamp;
}

public boolean isMultiplexedSession() {
return isMultiplexedSession;
}

@Override
public String toString() {
return String.format(
"transactionId: %s, sessionId: %s, timestamp: %s, isMultiplexedSession: %s",
transactionId.toStringUtf8(), sessionId, timestamp, isMultiplexedSession);
"transactionId: %s, sessionId: %s, timestamp: %s",
transactionId.toStringUtf8(), sessionId, timestamp);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -317,7 +317,7 @@ public long executePartitionedUpdate(final Statement stmt, final UpdateOption...
if (useMultiplexedSessionPartitionedOps) {
return getMultiplexedSession().executePartitionedUpdate(stmt, options);
}
return executePartitionedUpdateSession(stmt, options);
return executePartitionedUpdateWithPooledSession(stmt, options);
}

private long executePartitionedUpdateWithPooledSession(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -232,15 +232,9 @@ private SessionReference getSessionReference() {
*/
@Override
public long executePartitionedUpdate(Statement stmt, UpdateOption... options) {
try {
SessionReference sessionReference = getSessionReference();
return new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ true)
.executePartitionedUpdate(stmt, options);
} catch (InterruptedException e) {
throw new RuntimeException(e);
} catch (ExecutionException e) {
throw new RuntimeException(e);
}
SessionReference sessionReference = getSessionReference();
return new MultiplexedSessionTransaction(
client, span, sessionReference, NO_CHANNEL_HINT, /* singleUse = */ true)
.executePartitionedUpdate(stmt, options);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -556,7 +556,8 @@ public AsyncTransactionManager transactionManagerAsync(TransactionOption... opti

@Override
public long executePartitionedUpdate(Statement stmt, UpdateOption... options) {
return createMultiplexedSessionTransaction(/* singleUse = */ true).executePartitionedUpdate(stmt, options);
return createMultiplexedSessionTransaction(/* singleUse = */ true)
.executePartitionedUpdate(stmt, options);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -414,16 +414,10 @@ private List<SessionImpl> internalBatchCreateSessions(

/** Returns a {@link SessionImpl} that references the existing session with the given name. */
SessionImpl sessionWithId(String name) {
return sessionWithId(name, /*isMultiplexedSession= */ false);
}

/** Returns a {@link SessionImpl} that references the existing session with the given name. */
SessionImpl sessionWithId(String name, boolean isMultiplexedSession) {
final Map<SpannerRpc.Option, ?> options;
synchronized (this) {
options = optionMap(SessionOption.channelHint(sessionChannelCounter++));
}
return new SessionImpl(
spanner, new SessionReference(name, /*createTime= */ null, isMultiplexedSession, options));
return new SessionImpl(spanner, new SessionReference(name, options));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -132,14 +132,12 @@ public void testBatchReadOnlyTxnWithBound() throws Exception {
assertThat(batchTxn.getReadTimestamp()).isEqualTo(t);
assertThat(batchTxn.getReadTimestamp())
.isEqualTo(batchTxn.getBatchTransactionId().getTimestamp());
assertEquals(batchTxn.getBatchTransactionId().isMultiplexedSession(), isMultiplexedSession);
}

@Test
public void testBatchReadOnlyTxnWithTxnId() {
when(txnID.getSessionId()).thenReturn(SESSION_NAME);
when(txnID.getTransactionId()).thenReturn(TXN_ID);
when(txnID.isMultiplexedSession()).thenReturn(isMultiplexedSession);
Timestamp t = Timestamp.parseTimestamp(TIMESTAMP);
when(txnID.getTimestamp()).thenReturn(t);

Expand All @@ -149,8 +147,6 @@ public void testBatchReadOnlyTxnWithTxnId() {
assertThat(batchTxn.getReadTimestamp()).isEqualTo(t);
assertThat(batchTxn.getReadTimestamp())
.isEqualTo(batchTxn.getBatchTransactionId().getTimestamp());
assertThat(batchTxn.getBatchTransactionId().isMultiplexedSession())
.isEqualTo(isMultiplexedSession);
}

@Test
Expand Down