From 2e322a790ea04e8934ed63840787463e8c617b0c Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Sat, 3 Oct 2026 10:37:29 -0700 Subject: [PATCH] feat: expose backpressure and batching options in the perf client Add --max-pending-bytes, --max-write-batches-in-flight, --max-read-batches-in-flight and --batching-threads to the perf client, passed through to the corresponding OxiaClientBuilder options. --max-pending-bytes accepts sizes with a binary unit suffix (e.g. 512K, 10M, 1G). Signed-off-by: Matteo Merli --- .../io/oxia/client/perf/PerfArguments.java | 23 ++++++++ .../java/io/oxia/client/perf/PerfClient.java | 4 ++ .../io/oxia/client/perf/SizeConverter.java | 52 +++++++++++++++++++ .../oxia/client/perf/SizeConverterTest.java | 51 ++++++++++++++++++ 4 files changed, 130 insertions(+) create mode 100644 perf/src/main/java/io/oxia/client/perf/SizeConverter.java create mode 100644 perf/src/test/java/io/oxia/client/perf/SizeConverterTest.java diff --git a/perf/src/main/java/io/oxia/client/perf/PerfArguments.java b/perf/src/main/java/io/oxia/client/perf/PerfArguments.java index 99cc1342..beb39014 100644 --- a/perf/src/main/java/io/oxia/client/perf/PerfArguments.java +++ b/perf/src/main/java/io/oxia/client/perf/PerfArguments.java @@ -68,6 +68,29 @@ public class PerfArguments { description = "Requests timeout") long requestTimeoutMs = OxiaClientBuilderImpl.DefaultRequestTimeout.toMillis(); + @Parameter( + names = {"--max-pending-bytes"}, + description = + "Max total size of pending operations in the client, e.g. 512K, 10M, 1G" + + " (0 to disable)", + converter = SizeConverter.class) + long maxPendingBytes = OxiaClientBuilderImpl.DefaultMaxPendingBytes; + + @Parameter( + names = {"--max-write-batches-in-flight"}, + description = "Max number of in-flight write batches per shard") + int maxWriteBatchesInFlight = OxiaClientBuilderImpl.DefaultMaxWriteBatchesInFlight; + + @Parameter( + names = {"--max-read-batches-in-flight"}, + description = "Max number of in-flight read batches per shard") + int maxReadBatchesInFlight = OxiaClientBuilderImpl.DefaultMaxReadBatchesInFlight; + + @Parameter( + names = {"--batching-threads"}, + description = "Number of threads assembling operation batches") + int batchingThreads = OxiaClientBuilderImpl.DefaultBatchingThreads; + @Parameter( names = {"-o", "--max-outstanding-requests"}, description = "Max number of outstanding requests to server") diff --git a/perf/src/main/java/io/oxia/client/perf/PerfClient.java b/perf/src/main/java/io/oxia/client/perf/PerfClient.java index 8e70cbf2..d13b2f94 100644 --- a/perf/src/main/java/io/oxia/client/perf/PerfClient.java +++ b/perf/src/main/java/io/oxia/client/perf/PerfClient.java @@ -86,6 +86,10 @@ public static void main(String[] args) throws Exception { AsyncOxiaClient client = OxiaClientBuilder.create(arguments.serviceAddr) .maxRequestsPerBatch(arguments.maxRequestsPerBatch) + .maxPendingBytes(arguments.maxPendingBytes) + .maxWriteBatchesInFlight(arguments.maxWriteBatchesInFlight) + .maxReadBatchesInFlight(arguments.maxReadBatchesInFlight) + .batchingThreads(arguments.batchingThreads) .requestTimeout(Duration.ofMillis(arguments.requestTimeoutMs)) .namespace(arguments.namespace) .openTelemetry(sdk.getOpenTelemetrySdk()) diff --git a/perf/src/main/java/io/oxia/client/perf/SizeConverter.java b/perf/src/main/java/io/oxia/client/perf/SizeConverter.java new file mode 100644 index 00000000..9ce7111b --- /dev/null +++ b/perf/src/main/java/io/oxia/client/perf/SizeConverter.java @@ -0,0 +1,52 @@ +/* + * Copyright © 2022-2025 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.perf; + +import com.beust.jcommander.ParameterException; +import com.beust.jcommander.converters.BaseConverter; + +/** Parses a size in bytes, with an optional binary unit suffix: K, M, G or T (e.g. 10M, 1G). */ +public class SizeConverter extends BaseConverter { + + public SizeConverter(String optionName) { + super(optionName); + } + + @Override + public Long convert(String value) { + try { + return parseSize(value); + } catch (NumberFormatException | ArithmeticException e) { + throw new ParameterException(getErrorString(value, "a size (e.g. 512K, 10M, 1G)")); + } + } + + static long parseSize(String value) { + String s = value.trim(); + int shift = + switch (s.isEmpty() ? ' ' : Character.toUpperCase(s.charAt(s.length() - 1))) { + case 'K' -> 10; + case 'M' -> 20; + case 'G' -> 30; + case 'T' -> 40; + default -> 0; + }; + if (shift > 0) { + s = s.substring(0, s.length() - 1); + } + return Math.multiplyExact(Long.parseLong(s), 1L << shift); + } +} diff --git a/perf/src/test/java/io/oxia/client/perf/SizeConverterTest.java b/perf/src/test/java/io/oxia/client/perf/SizeConverterTest.java new file mode 100644 index 00000000..3a09150e --- /dev/null +++ b/perf/src/test/java/io/oxia/client/perf/SizeConverterTest.java @@ -0,0 +1,51 @@ +/* + * Copyright © 2022-2025 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.perf; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import com.beust.jcommander.ParameterException; +import org.junit.jupiter.api.Test; + +class SizeConverterTest { + + private final SizeConverter converter = new SizeConverter("--max-pending-bytes"); + + @Test + void plainBytes() { + assertThat(converter.convert("0")).isEqualTo(0L); + assertThat(converter.convert("1234")).isEqualTo(1234L); + } + + @Test + void unitSuffixes() { + assertThat(converter.convert("512K")).isEqualTo(512L * 1024); + assertThat(converter.convert("10M")).isEqualTo(10L * 1024 * 1024); + assertThat(converter.convert("1G")).isEqualTo(1024L * 1024 * 1024); + assertThat(converter.convert("2T")).isEqualTo(2L * 1024 * 1024 * 1024 * 1024); + assertThat(converter.convert("256m")).isEqualTo(256L * 1024 * 1024); + } + + @Test + void invalid() { + assertThatThrownBy(() -> converter.convert("")).isInstanceOf(ParameterException.class); + assertThatThrownBy(() -> converter.convert("M")).isInstanceOf(ParameterException.class); + assertThatThrownBy(() -> converter.convert("10X")).isInstanceOf(ParameterException.class); + assertThatThrownBy(() -> converter.convert("1.5G")).isInstanceOf(ParameterException.class); + assertThatThrownBy(() -> converter.convert("10000000T")).isInstanceOf(ParameterException.class); + } +}