From a0d4f41a38104a20b20fecfe7642b3029018d50c Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sun, 4 Oct 2026 05:26:12 +0300 Subject: [PATCH] [improve][test] Create the IoT gateways' producers concurrently in a random order Assisted-by: Claude Code (claude-opus-5-5) --- .../scenarios/docs/iot-telemetry.md | 18 ++- .../scenarios/iot-telemetry-base.yaml | 1 + .../tests/performance/tools/IotScenario.java | 14 ++- .../performance/tools/TelemetryProducer.java | 100 +++++++++++++++-- .../performance/tools/IotScenarioTest.java | 19 +++- .../tools/TelemetryProducerTest.java | 104 ++++++++++++++++++ 6 files changed, 243 insertions(+), 13 deletions(-) create mode 100644 tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/TelemetryProducerTest.java diff --git a/tests/performance/scenarios/docs/iot-telemetry.md b/tests/performance/scenarios/docs/iot-telemetry.md index 376225662a7de..d3b041e16cdcf 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -112,9 +112,20 @@ separate scenario with normal or deliberately short limits when measuring rollov or long-running storage behavior. BookKeeper entry-log flushing and disk-space checks remain enabled. Set `rate: 0` together with a positive `measurement.messages` to remove producer pacing. Set -`gateways.producer.precreate: true` to open every gateway/topic producer before throughput timing begins. The gateways' -summary, `gateways-summary.json`, reports `messagesPerSecond` only for the post-warmup measurement phase and retains -`wholeRunMessagesPerSecond` as startup and warmup context. +`gateways.producer.precreate: true` to open every gateway/topic producer before throughput timing begins. + +With `precreate: true`, the gateways create their producers in a random order, up to +`gateways.producer.precreateConcurrency` (32) at a time, as independent gateways connect. Each gateway's client opens a +connection for its lookup and then one to the topic's broker, and a broker assigns the connections that it accepts to +its I/O threads in turn: in this test topology, where the service URL and the topics' owner are the same broker, +gateways that connected one at a time put every data connection on every other I/O thread, so half of the broker's +I/O threads served the traffic. Concurrent connections spread the data connections over the I/O threads on average, +not evenly. Runs made before this setting existed created the producers one at a time in the order of the gateways; +`precreateConcurrency: 1` reproduces that, and a scenario without the setting gets 32. With `precreate: false`, the +gateways create each producer when they first send to its topic, one at a time, so the pattern remains. + +The gateways' summary, `gateways-summary.json`, reports `messagesPerSecond` only for the post-warmup measurement phase +and retains `wholeRunMessagesPerSecond` as startup and warmup context. ## Settings @@ -148,6 +159,7 @@ workloads: maxOutstanding: 20000 # messages in flight across the gateways batchingEnabled: true precreate: false # open every producer before the first message + precreateConcurrency: 32 # with precreate, producers created at the same time, in a random order env: # the gateways' container, which runs every gateway; from the memory configuration PULSAR_MEM: -Xms512m -Xmx512m -XX:MaxDirectMemorySize=256m -XX:+UseTransparentHugePages -XX:+AlwaysPreTouch topics: diff --git a/tests/performance/scenarios/iot-telemetry-base.yaml b/tests/performance/scenarios/iot-telemetry-base.yaml index 9cbc717a28413..ac286813e8e3b 100644 --- a/tests/performance/scenarios/iot-telemetry-base.yaml +++ b/tests/performance/scenarios/iot-telemetry-base.yaml @@ -53,6 +53,7 @@ workloads: maxOutstanding: 20000 batchingEnabled: true precreate: false + precreateConcurrency: 32 topics: count: 30 prefix: persistent://public/default/iot-telemetry- diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/IotScenario.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/IotScenario.java index b3e263f9daefd..a2642c2f47e95 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/IotScenario.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/IotScenario.java @@ -76,9 +76,20 @@ public record Gateways(int count, Producer producer, Map env) { * * @param maxOutstanding the most messages in flight across the gateways * @param precreate open every gateway's producer of every topic before the first message + * @param precreateConcurrency with {@code precreate}, the most producers that the gateways create at the same + * time, in a random order of the gateways and topics; 1 creates them one at a time in + * order */ public record Producer(int ioThreads, int listenerThreads, int maxOutstanding, boolean batchingEnabled, - boolean precreate) { + boolean precreate, Integer precreateConcurrency) { + /** The default of {@code precreateConcurrency}, for scenarios written before it existed. */ + public static final int DEFAULT_PRECREATE_CONCURRENCY = 32; + + public Producer { + if (precreateConcurrency == null) { + precreateConcurrency = DEFAULT_PRECREATE_CONCURRENCY; + } + } } /** The telemetry topics, named {@code }. */ @@ -176,6 +187,7 @@ boolean enabled() { && client.listenerThreads() >= 1, "The gateways' and the applications' ioThreads and listenerThreads must be at least 1"); require(producer.maxOutstanding() >= 1, "gateways.producer.maxOutstanding must be at least 1"); + require(producer.precreateConcurrency() >= 1, "gateways.producer.precreateConcurrency must be at least 1"); if (timeoutSeconds < minimumRuntimeSeconds) { long warmupTotalSeconds = minimumRuntimeSeconds - durationSeconds; throw new IllegalArgumentException(String.format(Locale.ROOT, "Invalid IoT scenario: timeoutSeconds is %d," diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryProducer.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryProducer.java index a8f5966fef798..b922221e40596 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryProducer.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryProducer.java @@ -23,17 +23,23 @@ import java.nio.ByteBuffer; import java.nio.file.Files; import java.util.ArrayList; +import java.util.Collections; import java.util.List; +import java.util.Random; import java.util.Set; import java.util.SplittableRandom; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.LockSupport; +import java.util.function.IntFunction; import org.apache.pulsar.client.api.BatcherBuilder; import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.ProducerBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientSharedResources; import picocli.CommandLine.Command; @@ -102,12 +108,7 @@ public Integer call() throws Exception { .build()); } if (scenario.gateways().producer().precreate()) { - for (int gateway = 0; gateway < scenario.gatewayCount(); gateway++) { - for (int topic = 0; topic < scenario.topicCount(); topic++) { - int producerIndex = gateway * scenario.topicCount() + topic; - producers[producerIndex] = createProducer(scenario, clients, gateway, topic); - } - } + precreateProducers(scenario, clients, producers); } phase = scenario.warmupMessageCount() > 0 ? "warmup" : "measurement"; SplittableRandom random = new SplittableRandom(0x51c0ffeeL); @@ -267,14 +268,97 @@ private static void awaitOutstanding(Semaphore outstanding, int permits) throws private Producer createProducer(IotScenario scenario, List clients, int gateway, int topic) throws Exception { + return producerBuilder(scenario, clients, gateway, topic).create(); + } + + /** + * Creates every gateway's producer of every topic, up to {@code precreateConcurrency} at a time, in a random + * order. Real gateways connect independently, many at the same time, for example when they reconnect after a + * broker restart. A gateway's client connects to the service URL for its lookup and then to the topic's broker; + * created one at a time, every gateway's two connections would reach the broker one after the other, and a broker, + * which assigns the connections that it accepts to its I/O threads in turn, would serve every data connection on + * every other I/O thread. After a creation fails, no more are started, and the failure is thrown once the started + * ones have completed. + */ + private void precreateProducers(IotScenario scenario, List clients, Producer[] producers) + throws Exception { + int topicCount = scenario.topicCount(); + int concurrency = scenario.gateways().producer().precreateConcurrency(); + precreate(precreateOrder(producers.length, concurrency), concurrency, producers, + producerIndex -> producerBuilder(scenario, clients, producerIndex / topicCount, + producerIndex % topicCount).createAsync()); + } + + /** + * Creates the objects in the given order, up to {@code concurrency} at a time, into {@code created} at their + * index. After a creation fails, no more are started, and the first failure is thrown once the started ones have + * completed. + */ + static void precreate(int[] order, int concurrency, T[] created, IntFunction> create) + throws Exception { + Semaphore permits = new Semaphore(concurrency); + AtomicReference firstFailure = new AtomicReference<>(); + List> creations = new ArrayList<>(order.length); + for (int index : order) { + permits.acquire(); + if (firstFailure.get() != null) { + permits.release(); + break; + } + CompletableFuture creation; + try { + creation = create.apply(index); + } catch (RuntimeException e) { + permits.release(); + firstFailure.compareAndSet(null, e); + break; + } + creations.add(creation.whenComplete((object, failure) -> { + if (object != null) { + created[index] = object; + } else { + firstFailure.compareAndSet(null, failure instanceof CompletionException && failure.getCause() + != null ? failure.getCause() : failure); + } + permits.release(); + })); + } + // Waiting for the started creations also makes the created objects visible to this thread, so that they're + // closed with the others if a creation failed + CompletableFuture.allOf(creations.toArray(new CompletableFuture[0])).handle((ignored, failure) -> null) + .get(); + Throwable failure = firstFailure.get(); + if (failure instanceof Exception exception) { + throw exception; + } else if (failure != null) { + throw new IllegalStateException("A creation failed", failure); + } + } + + /** + * The order in which the gateways' producers are precreated: the order of the gateways and topics for a + * concurrency of 1, otherwise a random order, the same for every run of a scenario. + */ + static int[] precreateOrder(int producerCount, int concurrency) { + List order = new ArrayList<>(producerCount); + for (int producerIndex = 0; producerIndex < producerCount; producerIndex++) { + order.add(producerIndex); + } + if (concurrency > 1) { + Collections.shuffle(order, new Random(0x9a7e3aL)); + } + return order.stream().mapToInt(Integer::intValue).toArray(); + } + + private ProducerBuilder producerBuilder(IotScenario scenario, List clients, + int gateway, int topic) { return clients.get(gateway).newProducer() .topic(scenario.topicNames().get(topic)) .producerName("iot-gateway-" + gateway + "-topic-" + topic) .batcherBuilder(BatcherBuilder.KEY_BASED) .enableBatching(scenario.gateways().producer().batchingEnabled()) .blockIfQueueFull(true) - .sendTimeout(0, TimeUnit.SECONDS) - .create(); + .sendTimeout(0, TimeUnit.SECONDS); } private void writeState(long[] sequences) throws Exception { diff --git a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/IotScenarioTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/IotScenarioTest.java index 4f2ae2955cd07..8d563c12bb3c0 100644 --- a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/IotScenarioTest.java +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/IotScenarioTest.java @@ -103,6 +103,23 @@ public void rejectsTimeBasedWarmupWithoutRateLimit() { .isInstanceOf(IllegalArgumentException.class); } + @Test + public void defaultsAnOmittedPrecreateConcurrency() { + assertThat(new IotScenario.Producer(2, 2, 100, true, true, null).precreateConcurrency()) + .isEqualTo(IotScenario.Producer.DEFAULT_PRECREATE_CONCURRENCY); + } + + @Test + public void rejectsPrecreateConcurrencyBelowOne() { + assertThatThrownBy(() -> new IotScenario("pulsar://localhost:6650", new IotScenario.Warmup(0, 0, 1, 0), + new IotScenario.Measurement(120, 1_000), 0, new IotScenario.Payload(64), new IotScenario.Devices(1_000), + new IotScenario.Gateways(10, new IotScenario.Producer(2, 2, 100, true, true, 0), null), + new IotScenario.Topics(2, "persistent://public/default/iot-"), + new IotScenario.Applications(1, 2, "app-", new IotScenario.Client(2, 2), null), null, 300)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("gateways.producer.precreateConcurrency must be at least 1"); + } + @Test public void rejectsTwoWarmupLimits() { assertThatThrownBy(() -> scenario(20, 1_000, 1000, 0)) @@ -136,7 +153,7 @@ private static IotScenario scenario(int warmupSeconds, long warmupMessages, int new IotScenario.Warmup(warmupSeconds, warmupMessages, warmupRounds, warmupRoundDelaySeconds), new IotScenario.Measurement(120, numberOfMessages), rate, new IotScenario.Payload(64), new IotScenario.Devices(1_000), - new IotScenario.Gateways(10, new IotScenario.Producer(2, 2, 100, true, true), null), + new IotScenario.Gateways(10, new IotScenario.Producer(2, 2, 100, true, true, 32), null), new IotScenario.Topics(2, "persistent://public/default/iot-"), new IotScenario.Applications(1, 2, "app-", new IotScenario.Client(2, 2), null), null, timeoutSeconds); } diff --git a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/TelemetryProducerTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/TelemetryProducerTest.java new file mode 100644 index 0000000000000..c457b226769ad --- /dev/null +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/TelemetryProducerTest.java @@ -0,0 +1,104 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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 org.apache.pulsar.tests.performance.tools; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Supplier; +import java.util.stream.IntStream; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +public class TelemetryProducerTest { + private ScheduledExecutorService executor; + + @BeforeMethod + public void startExecutor() { + executor = Executors.newScheduledThreadPool(4); + } + + @AfterMethod(alwaysRun = true) + public void stopExecutor() { + executor.shutdownNow(); + } + + @Test + public void precreateOrderIsTheGatewaysOrderOneAtATimeAndARepeatableRandomOrderConcurrently() { + assertThat(TelemetryProducer.precreateOrder(6, 1)).containsExactly(0, 1, 2, 3, 4, 5); + + int[] concurrent = TelemetryProducer.precreateOrder(500, 32); + // every producer once, not in the order of the gateways, and the same order for every run + assertThat(concurrent).containsExactlyInAnyOrder(IntStream.range(0, 500).toArray()); + assertThat(concurrent).isNotEqualTo(IntStream.range(0, 500).toArray()); + assertThat(TelemetryProducer.precreateOrder(500, 32)).isEqualTo(concurrent); + } + + @Test + public void precreatesUpToTheConcurrencyAtATime() throws Exception { + AtomicInteger inFlight = new AtomicInteger(); + AtomicInteger maxInFlight = new AtomicInteger(); + Integer[] created = new Integer[100]; + TelemetryProducer.precreate(IntStream.range(0, 100).toArray(), 8, created, index -> { + maxInFlight.accumulateAndGet(inFlight.incrementAndGet(), Math::max); + return completeLater(() -> { + inFlight.decrementAndGet(); + return index; + }); + }); + assertThat(created).containsExactly(IntStream.range(0, 100).boxed().toArray(Integer[]::new)); + assertThat(maxInFlight.get()).isBetween(1, 8); + } + + @Test + public void stopsAfterAFailureAndThrowsItOnceTheStartedCreationsHaveCompleted() { + IllegalStateException failure = new IllegalStateException("lookup failed"); + List started = new CopyOnWriteArrayList<>(); + List completed = new CopyOnWriteArrayList<>(); + Integer[] created = new Integer[100]; + assertThatThrownBy(() -> TelemetryProducer.precreate(IntStream.range(0, 100).toArray(), 4, created, index -> { + started.add(index); + if (index == 10) { + return CompletableFuture.failedFuture(failure); + } + return completeLater(() -> { + completed.add(index); + return index; + }); + })).isSameAs(failure); + // no more creations started after the failure was seen, and every started one completed before the throw + assertThat(started).hasSizeLessThan(100); + assertThat(started.size()).isLessThanOrEqualTo(10 + 1 + 4); + assertThat(completed).containsExactlyInAnyOrderElementsOf(started.stream().filter(i -> i != 10).toList()); + } + + // completes the value's future after a millisecond on another thread, as a producer's creation completes + private CompletableFuture completeLater(Supplier value) { + CompletableFuture future = new CompletableFuture<>(); + executor.schedule(() -> future.complete(value.get()), 1, TimeUnit.MILLISECONDS); + return future; + } +}