Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
@@ -0,0 +1,82 @@
/*
* Copyright © 2026 The Oxia Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.oxia.client.util;

import java.time.Duration;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.OutputTimeUnit;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
import org.openjdk.jmh.annotations.TearDown;
import org.openjdk.jmh.annotations.Threads;
import org.openjdk.jmh.annotations.Warmup;

/**
* Compares {@link TimeoutSweeper} against the {@link CompletableFuture#orTimeout} it replaced, for
* an operation that completes well before its timeout. Run with {@code ./gradlew :benchmarks:jmh};
* the threads contend on the JVM-wide delayer with {@code orTimeout}.
*/
@BenchmarkMode(Mode.Throughput)
@OutputTimeUnit(TimeUnit.MICROSECONDS)
@State(Scope.Benchmark)
@Fork(1)
@Threads(8)
@Warmup(iterations = 3, time = 1)
@Measurement(iterations = 5, time = 2)
public class TimeoutSweeperBenchmark {

private static final Duration TIMEOUT = Duration.ofSeconds(30);
private static final Object RESULT = new Object();

private ScheduledExecutorService executor;
private TimeoutSweeper sweeper;

@Setup
public void setup() {
executor = Executors.newSingleThreadScheduledExecutor();
sweeper = new TimeoutSweeper(executor, TIMEOUT);
}

@TearDown
public void tearDown() {
sweeper.close();
executor.shutdownNow();
}

@Benchmark
public CompletableFuture<Object> orTimeout() {
var future =
new CompletableFuture<Object>().orTimeout(TIMEOUT.toMillis(), TimeUnit.MILLISECONDS);
future.complete(RESULT);
return future;
}

@Benchmark
public CompletableFuture<Object> timeoutSweeper() {
var future = sweeper.add(new CompletableFuture<>());
future.complete(RESULT);
return future;
}
}
27 changes: 14 additions & 13 deletions client/src/main/java/io/oxia/client/AsyncOxiaClientImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
import io.oxia.client.session.SessionManager;
import io.oxia.client.shard.ShardManager;
import io.oxia.client.util.PendingBytesLimiter;
import io.oxia.client.util.TimeoutSweeper;
import io.oxia.proto.KeyComparisonType;
import io.oxia.proto.ListRequest;
import io.oxia.proto.ListResponse;
Expand All @@ -67,7 +68,6 @@
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
Expand Down Expand Up @@ -184,7 +184,7 @@ class AsyncOxiaClientImpl implements AsyncOxiaClient {
private final @NonNull BatchManager readBatchManager;
private final @NonNull BatchManager writeBatchManager;
private final @NonNull SessionManager sessionManager;
private final long requestTimeoutMs;
private final @NonNull TimeoutSweeper requestTimeouts;
private final @NonNull PendingBytesLimiter pendingBytesLimiter;
private volatile boolean closed;

Expand Down Expand Up @@ -242,7 +242,7 @@ class AsyncOxiaClientImpl implements AsyncOxiaClient {
this.sessionManager = sessionManager;
this.scheduledExecutor = scheduledExecutor;
this.ownsResources = ownsResources;
this.requestTimeoutMs = requestTimeout.toMillis();
this.requestTimeouts = new TimeoutSweeper(scheduledExecutor, requestTimeout);

counterPutBytes =
instrumentProvider.newCounter(
Expand Down Expand Up @@ -375,8 +375,8 @@ class AsyncOxiaClientImpl implements AsyncOxiaClient {
callback = CompletableFuture.failedFuture(e);
}
final long pendingBytes = acquiredBytes;
return callback
.orTimeout(requestTimeoutMs, TimeUnit.MILLISECONDS)
return requestTimeouts
.add(callback)
.whenComplete(
(putResult, throwable) -> {
if (pendingBytes > 0) {
Expand Down Expand Up @@ -488,8 +488,8 @@ private CompletableFuture<PutResult> internalPut(
callback.completeExceptionally(e);
}
final long pendingBytes = acquiredBytes;
return callback
.orTimeout(requestTimeoutMs, TimeUnit.MILLISECONDS)
return requestTimeouts
.add(callback)
.whenComplete(
(putResult, throwable) -> {
if (pendingBytes > 0) {
Expand Down Expand Up @@ -552,8 +552,8 @@ private CompletableFuture<PutResult> internalPut(
callback = CompletableFuture.failedFuture(e);
}
final long pendingBytes = acquiredBytes;
return callback
.orTimeout(requestTimeoutMs, TimeUnit.MILLISECONDS)
return requestTimeouts
.add(callback)
.whenComplete(
(putResult, throwable) -> {
if (pendingBytes > 0) {
Expand Down Expand Up @@ -593,8 +593,8 @@ private CompletableFuture<PutResult> internalPut(
callback.completeExceptionally(e);
}
final long pendingBytes = acquiredBytes;
return callback
.orTimeout(requestTimeoutMs, TimeUnit.MILLISECONDS)
return requestTimeouts
.add(callback)
.whenComplete(
(getResult, throwable) -> {
if (pendingBytes > 0) {
Expand Down Expand Up @@ -697,8 +697,8 @@ private void internalGetMultiShards(
} catch (Exception e) {
callback = CompletableFuture.failedFuture(e);
}
return callback
.orTimeout(requestTimeoutMs, TimeUnit.MILLISECONDS)
return requestTimeouts
.add(callback)
.whenComplete(
(listResult, throwable) -> {
gaugePendingListRequests.decrement();
Expand Down Expand Up @@ -965,6 +965,7 @@ public void close() throws Exception {
// In shared mode the RpcProvider does not own the connection pool, so this only closes the
// per-client write streams; the shared connections stay open for other clients.
rpcProvider.close();
requestTimeouts.close();
if (ownsResources) {
scheduledExecutor.shutdownNow();
}
Expand Down
23 changes: 9 additions & 14 deletions client/src/main/java/io/oxia/client/grpc/GrpcRpcProvider.java
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import io.oxia.client.ClientConfig;
import io.oxia.client.grpc.observer.CancelableStreamObserver;
import io.oxia.client.grpc.observer.ManagedObservers;
import io.oxia.client.util.TimeoutSweeper;
import io.oxia.proto.CloseSessionRequest;
import io.oxia.proto.CloseSessionResponse;
import io.oxia.proto.CreateSessionRequest;
Expand Down Expand Up @@ -71,6 +72,7 @@ final class GrpcRpcProvider implements RpcProvider {
private final ScheduledExecutorService asyncExecutor;
private final LongFunction<String> shardLeaderProvider;
private final Map<Long, ManagedWriteStream> writeStreams;
private final TimeoutSweeper requestTimeouts;

GrpcRpcProvider(
@NonNull ClientConfig clientConfig,
Expand Down Expand Up @@ -104,6 +106,7 @@ private GrpcRpcProvider(
this.ownsConnectionManager = ownsConnectionManager;
this.shardLeaderProvider = shardLeaderProvider;
this.writeStreams = Maps.newConcurrentMap();
this.requestTimeouts = new TimeoutSweeper(asyncExecutor, clientConfig.requestTimeout());
}

@Override
Expand All @@ -117,9 +120,7 @@ public void getShardAssignments(
.with(asyncExecutor)
.getStageAsync(
() -> {
final var barrierFuture =
new CompletableFuture<Void>()
.orTimeout(clientConfig.requestTimeout().toMillis(), TimeUnit.MILLISECONDS);
final var barrierFuture = requestTimeouts.add(new CompletableFuture<Void>());
final var barrierObserver =
ManagedObservers.toBarrierStreamObserver(guardedObserver, barrierFuture);
final var attemptContext = Context.current().withCancellation();
Expand Down Expand Up @@ -165,8 +166,7 @@ public void getNotifications(
// resumed one gets nothing until a new notification is written, so it is only
// bounded by the subscription max age, like the sequence updates.
if (!request.hasStartOffsetExclusive()) {
barrierFuture.orTimeout(
clientConfig.requestTimeout().toMillis(), TimeUnit.MILLISECONDS);
requestTimeouts.add(barrierFuture);
}
final var barrierObserver =
ManagedObservers.toBarrierClientResponseObserver(observer, barrierFuture);
Expand Down Expand Up @@ -284,9 +284,7 @@ public void read(@NonNull ReadRequest request, @NonNull StreamObserver<ReadRespo
.with(asyncExecutor)
.getStageAsync(
() -> {
final var barrierFuture =
new CompletableFuture<Void>()
.orTimeout(clientConfig.requestTimeout().toMillis(), TimeUnit.MILLISECONDS);
final var barrierFuture = requestTimeouts.add(new CompletableFuture<Void>());
final var barrierObserver =
ManagedObservers.toBarrierStreamObserver(guardedObserver, barrierFuture);
final var attemptContext = Context.current().withCancellation();
Expand Down Expand Up @@ -364,9 +362,7 @@ public void list(
.with(asyncExecutor)
.getStageAsync(
() -> {
final var barrierFuture =
new CompletableFuture<Void>()
.orTimeout(clientConfig.requestTimeout().toMillis(), TimeUnit.MILLISECONDS);
final var barrierFuture = requestTimeouts.add(new CompletableFuture<Void>());
final var barrierObserver =
ManagedObservers.toBarrierClientResponseObserver(observer, barrierFuture);
final var attemptContext = Context.current().withCancellation();
Expand Down Expand Up @@ -410,9 +406,7 @@ public void rangeScan(
.with(asyncExecutor)
.getStageAsync(
() -> {
final var barrierFuture =
new CompletableFuture<Void>()
.orTimeout(clientConfig.requestTimeout().toMillis(), TimeUnit.MILLISECONDS);
final var barrierFuture = requestTimeouts.add(new CompletableFuture<Void>());
final var barrierObserver =
ManagedObservers.toBarrierClientResponseObserver(observer, barrierFuture);
final var attemptContext = Context.current().withCancellation();
Expand Down Expand Up @@ -495,6 +489,7 @@ private OxiaClientGrpc.OxiaClientStub withSubscriptionMaxAge(OxiaClientGrpc.Oxia

@Override
public void close() throws Exception {
requestTimeouts.close();
try {
writeStreams.values().forEach(ManagedWriteStream::close);
writeStreams.clear();
Expand Down
Loading
Loading