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
18 changes: 15 additions & 3 deletions tests/performance/scenarios/docs/iot-telemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -128,9 +128,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

Expand Down Expand Up @@ -164,6 +175,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:
Expand Down
1 change: 1 addition & 0 deletions tests/performance/scenarios/iot-telemetry-base.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ workloads:
maxOutstanding: 20000
batchingEnabled: true
precreate: false
precreateConcurrency: 32
topics:
count: 30
prefix: persistent://public/default/iot-telemetry-
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,9 +77,20 @@ public record Gateways(int count, Producer producer, Map<String, String> 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 <prefix><index>}. */
Expand Down Expand Up @@ -207,6 +218,7 @@ boolean enabled() {
"applications.joinSeconds must be at least 0 with no null value, and the warmup ("
+ warmupTotalSeconds + " s) and the latest join must be within timeoutSeconds");
require(applications.caughtUpLatencyMillis() >= 1, "applications.caughtUpLatencyMillis must be at least 1");
require(producer.precreateConcurrency() >= 1, "gateways.producer.precreateConcurrency must be at least 1");
if (timeoutSeconds < minimumRuntimeSeconds) {
throw new IllegalArgumentException(String.format(Locale.ROOT, "Invalid IoT scenario: timeoutSeconds is %d,"
+ " but the workload needs %d s: %d s of warmup (%d round(s) of %d s%s) and %d s of"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -271,14 +272,97 @@ private static void awaitOutstanding(Semaphore outstanding, int permits) throws

private Producer<byte[]> createProducer(IotScenario scenario, List<PulsarClient> 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<PulsarClient> clients, Producer<byte[]>[] 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 <T> void precreate(int[] order, int concurrency, T[] created, IntFunction<CompletableFuture<T>> create)
throws Exception {
Semaphore permits = new Semaphore(concurrency);
AtomicReference<Throwable> firstFailure = new AtomicReference<>();
List<CompletableFuture<?>> creations = new ArrayList<>(order.length);
for (int index : order) {
permits.acquire();
if (firstFailure.get() != null) {
permits.release();
break;
}
CompletableFuture<T> 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<Integer> 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<byte[]> producerBuilder(IotScenario scenario, List<PulsarClient> 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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,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))
Expand Down Expand Up @@ -197,7 +214,7 @@ private static IotScenario scenario(IotScenario.Applications applications, IotSc
// a time-based warmup needs a rate
return new IotScenario("pulsar://localhost:6650", warmup, new IotScenario.Measurement(120, 1_000),
warmup.seconds() > 0 ? 100 : 0, 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, null), null),
new IotScenario.Topics(2, "persistent://public/default/iot-"), applications, null, 300);
}

Expand All @@ -208,7 +225,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);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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<Integer> started = new CopyOnWriteArrayList<>();
List<Integer> 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 <T> CompletableFuture<T> completeLater(Supplier<T> value) {
CompletableFuture<T> future = new CompletableFuture<>();
executor.schedule(() -> future.complete(value.get()), 1, TimeUnit.MILLISECONDS);
return future;
}
}
Loading