From 693ce8150cc7bf429d33329b35028766d7e59951 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sat, 3 Oct 2026 22:05:45 +0300 Subject: [PATCH 1/5] [improve][test] Add a catch-up scenario: applications joining a live stream after the measurement starts applications.joinSeconds has applications join after the measurement starts. A late application creates its subscription at the earliest position before the gateways publish, so that its backlog builds up, skips the warmup, and opens its pods when the gateways' measurement start, which they write to the coordination directory, plus its join time has passed. It has caught up when it first receives a measured message within applications.caughtUpLatencyMillis of its publishing. The run report's Catch-up section shows each late application's join time, its subscription's backlog then, how long it took to catch up and its catch-up rate. iot-telemetry-catch-up.yaml publishes the high-rate scenario's 30,000 msg/s for 120 s, with 4 of its 5 applications joining at 20 s (2 of them), 40 s and 60 s. Assisted-by: Claude Code (claude-opus-5-5) --- tests/performance/docs/run-reports.md | 3 + .../tests/performance/report/RunReport.java | 53 +++++++- .../performance/report/RunReportTest.java | 32 +++++ tests/performance/scenarios/README.md | 1 + .../scenarios/docs/iot-telemetry.md | 17 +++ .../scenarios/iot-telemetry-catch-up.yaml | 31 +++++ .../tests/performance/tools/IotScenario.java | 38 +++++- .../tools/MeasurementStartMarker.java | 62 ++++++++++ .../performance/tools/PerformanceTool.java | 4 + .../performance/tools/TelemetryConsumer.java | 115 +++++++++++++++++- .../performance/tools/TelemetryProducer.java | 4 + .../performance/tools/IotScenarioTest.java | 34 ++++++ .../tools/MeasurementStartMarkerTest.java | 50 ++++++++ 13 files changed, 436 insertions(+), 8 deletions(-) create mode 100644 tests/performance/scenarios/iot-telemetry-catch-up.yaml create mode 100644 tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java create mode 100644 tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java diff --git a/tests/performance/docs/run-reports.md b/tests/performance/docs/run-reports.md index 59a5a090d093a..89366f1edd34a 100644 --- a/tests/performance/docs/run-reports.md +++ b/tests/performance/docs/run-reports.md @@ -143,6 +143,9 @@ report has these sections: runs can be compared. - **Throughput**: the gateways' throughput, the delivered throughput until the slowest application received the last message, the measurement's duration and how long the applications were still receiving after the gateways finished. +- **Catch-up**, when applications join after the measurement starts (`applications.joinSeconds`): each late + application's join time, its subscription's backlog then, how long it took to catch up, and its catch-up rate. The + backlog chart shows each subscription's backlog building up until its application joins and falling as it catches up. - **Latency**: the publish latency (send to acknowledgment) and each application's end-to-end latency (publish to consume) at percentiles from p50 to the maximum, with charts by percentile and over time. - **Backlog and rates**: each subscription's backlog and the per-second rates, sampled from the broker's topic diff --git a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java index 7d611339e98c1..666fcd28359da 100644 --- a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java +++ b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java @@ -322,10 +322,12 @@ public static Path write(Path runDirectory, Run run, ObjectMapper mapper) throws appendProfiles(report, runDirectory, mapper); appendCorrectness(report, run.workload(), consumers); appendThroughput(report, producer, consumers); - appendLatency(report, runDirectory, run, consumerHistograms, measurementStart); Path stats = runDirectory.resolve(TOPIC_STATS_FILE); - if (Files.isRegularFile(stats)) { - appendTopicStats(report, runDirectory, readSamples(stats), measurementStart, measurementEnd, + Samples samples = Files.isRegularFile(stats) ? readSamples(stats) : null; + appendCatchUp(report, run.workload(), consumers, samples, measurementStart); + appendLatency(report, runDirectory, run, consumerHistograms, measurementStart); + if (samples != null) { + appendTopicStats(report, runDirectory, samples, measurementStart, measurementEnd, run.cooldowns(), chartFooter(run.info(), run.finished())); } if (hostSamples != null) { @@ -642,6 +644,51 @@ private static void appendThroughput(StringBuilder report, JsonNode producer, Li String.format(Locale.ROOT, "%.1f s", Math.max(0, lastReceived - end) / 1000.0))); } + /** + * The applications that joined after the measurement started ({@code applications.joinSeconds}): when each + * joined, its subscription's backlog then, from the sampled topic stats, and how long it took to catch up. + */ + static void appendCatchUp(StringBuilder report, JsonNode workload, List consumers, Samples samples, + long measurementStart) { + List late = consumers.stream().filter(consumer -> consumer.path("joinEpochMs").asLong() > 0) + .toList(); + if (late.isEmpty()) { + return; + } + int caughtUpLatency = workload.path("applications").path("caughtUpLatencyMillis").asInt(1000); + report.append("\n## Catch-up\n\nThe applications that joined after the measurement started. Each had its" + + " subscription from the start, so its backlog built up until it joined. An application has" + + " caught up when it first received a measured message within ") + .append(String.format(Locale.ROOT, "%,d", caughtUpLatency)) + .append(" ms of its publishing; its catch-up rate is the messages that it received until then, per" + + " second since it joined. The backlog is the sampled topic stats' last sample before the" + + " application joined.\n\n| Application | Joined | Backlog when it joined | Caught up after" + + " | Messages received until then | Catch-up rate |\n|---|---:|---:|---:|---:|---:|\n"); + for (JsonNode consumer : late) { + String application = applicationName(workload, consumer.path("applicationIndex").asInt()); + long joined = consumer.path("joinEpochMs").asLong(); + long caughtUp = consumer.path("caughtUpEpochMs").asLong(); + long messages = consumer.path("messagesWhenCaughtUp").asLong(); + String backlog = "–"; + double[] backlogs = samples != null ? samples.backlog().get(application) : null; + if (backlogs != null) { + for (int round = samples.epochMillis().length - 1; round >= 0; round--) { + if (samples.epochMillis()[round] <= joined && !Double.isNaN(backlogs[round])) { + backlog = String.format(Locale.ROOT, "%,.0f", backlogs[round]); + break; + } + } + } + report.append(String.format(Locale.ROOT, "| `%s` | %.1f s | %s | %s | %s | %s |%n", application, + (joined - measurementStart) / 1000.0, backlog, + caughtUp > 0 ? String.format(Locale.ROOT, "%.1f s", (caughtUp - joined) / 1000.0) + : "not caught up", + caughtUp > 0 ? String.format(Locale.ROOT, "%,d", messages) : "–", + caughtUp > joined ? String.format(Locale.ROOT, "%,.0f msg/s", + messages * 1000.0 / (caughtUp - joined)) : "–")); + } + } + /** * An application's name: its subscription, the workload's {@code applications.subscriptionPrefix} and its index, * such as {@code iot-application-0}, as the throughput and backlog charts name it. Each application consumes diff --git a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java index 4705fd4f3e672..323d92f9bcc94 100644 --- a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java +++ b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java @@ -28,6 +28,7 @@ import java.time.ZonedDateTime; import java.util.Comparator; import java.util.List; +import java.util.Map; import java.util.stream.Stream; import org.HdrHistogram.Histogram; import org.HdrHistogram.HistogramLogWriter; @@ -62,6 +63,37 @@ public void deleteRun() throws IOException { } } + @Test + public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception { + JsonNode workload = mapper.readTree("{\"applications\": {\"subscriptionPrefix\": \"app-\"," + + " \"caughtUpLatencyMillis\": 500}}"); + long start = 1_000_000; + List consumers = List.of( + mapper.readTree("{\"applicationIndex\": 0, \"joinEpochMs\": 0}"), + mapper.readTree("{\"applicationIndex\": 1, \"joinEpochMs\": " + (start + 20_000) + + ", \"caughtUpEpochMs\": " + (start + 30_000) + ", \"messagesWhenCaughtUp\": 1500000}"), + mapper.readTree("{\"applicationIndex\": 2, \"joinEpochMs\": " + (start + 40_000) + + ", \"caughtUpEpochMs\": 0, \"messagesWhenCaughtUp\": 0}")); + long[] epochs = {start + 19_000, start + 20_500, start + 39_000}; + RunReport.Samples samples = new RunReport.Samples(epochs, new double[3], Map.of(), + Map.of("app-1", new double[] {570_000, 600_000, 0}, "app-2", new double[] {1_100_000, 1_150_000, + 1_170_000})); + StringBuilder report = new StringBuilder(); + RunReport.appendCatchUp(report, workload, consumers, samples, start); + + assertThat(report.toString()) + .contains("## Catch-up") + .contains("within 500 ms of its publishing") + // the backlog of the last sample before the application joined, and its catch-up rate + .contains("| `app-1` | 20.0 s | 570,000 | 10.0 s | 1,500,000 | 150,000 msg/s |") + .contains("| `app-2` | 40.0 s | 1,170,000 | not caught up | – | – |") + .doesNotContain("`app-0`"); + // without late applications, there's no section + StringBuilder none = new StringBuilder(); + RunReport.appendCatchUp(none, workload, consumers.subList(0, 1), samples, start); + assertThat(none.toString()).isEmpty(); + } + @Test public void reportsCorrectnessThroughputLatencyAndSampledStats() throws IOException { Files.createDirectories(run.resolve("gateways")); diff --git a/tests/performance/scenarios/README.md b/tests/performance/scenarios/README.md index 35d0fbc32a355..0eec53c074510 100644 --- a/tests/performance/scenarios/README.md +++ b/tests/performance/scenarios/README.md @@ -47,6 +47,7 @@ detail. | [`iot-telemetry-small-restarts.yaml`](iot-telemetry-small-restarts.yaml) | The smaller topology, restarting 10 % of each application's clients every 30 s | low, about 3 GB | | [`iot-telemetry-restarts.yaml`](iot-telemetry-restarts.yaml) | The full topology, restarting 10 % of each application's clients every 30 s | medium, about 11 GB | | [`iot-telemetry-high-rate.yaml`](iot-telemetry-high-rate.yaml) | A high rate: five million messages at 30,000 messages per second, from 500 preconnected producers to one topic, consumed by 5 applications with 10 pods each | high, about 14 GB | +| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 120 s, with 4 of its 5 applications joining 20 s (2 of them), 40 s and 60 s after the measurement starts and catching up on their backlogs | high, about 14 GB | | [`iot-telemetry-max-rate.yaml`](iot-telemetry-max-rate.yaml) | The maximum rate: the 500 gateways of the high-rate scenario publish four million unbatched messages to one topic without a rate limit, consumed by one application with 20 pods, on single-copy ledgers | high, about 14 GB | ## Configurations diff --git a/tests/performance/scenarios/docs/iot-telemetry.md b/tests/performance/scenarios/docs/iot-telemetry.md index 376225662a7de..85046f7584f76 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -43,6 +43,13 @@ application-visible order across Key_Shared hash-range reassignment. workstation whose CPU runs at a fixed base frequency: without a rate limit, the gateways publish faster than the applications receive, and an application that falls behind can stall for tens of seconds. It uses the high-memory configuration. +- [`iot-telemetry-catch-up.yaml`](../iot-telemetry-catch-up.yaml) measures how quickly consumers that join a live + stream catch up. The gateways publish at the high-rate scenario's 30,000 messages per second for 120 s, and 4 of + the 5 applications join after the measurement starts: 2 at 20 s, and 1 each at 40 s and 60 s. Each late + application has its subscription from the start, so its backlog builds up until it joins, and it then reads the + backlog while the gateways keep publishing. The report's Catch-up section shows each application's backlog when it + joined and how long it took to catch up. The 2 applications that join together read the same backlog, so the + entries that one of them reads from storage can be served to the other from the broker's entry cache. - [`iot-telemetry-max-rate.yaml`](../iot-telemetry-max-rate.yaml) runs the high-rate scenario's 500 gateways at the maximum rate: `rate: 0` removes the rate limit, so the gateways publish four million unbatched 128-byte messages as fast as the cluster takes them, with at most 100,000 in flight. One application consumes them with twenty pods @@ -162,12 +169,22 @@ workloads: listenerThreads: 16 env: # the applications' container, which runs every application; from the memory configuration PULSAR_MEM: -Xms1536m -Xmx1536m -XX:MaxDirectMemorySize=256m -XX:+UseTransparentHugePages -XX:+AlwaysPreTouch + joinSeconds: [] # each application's join, in seconds after the measurement starts; empty: all at 0 + caughtUpLatencyMillis: 1000 # a late application has caught up within this end-to-end latency behaviors: podRestarts: # each application restarts this fraction of its pods every intervalSeconds intervalSeconds: 0 fraction: 0.0 ``` +`applications.joinSeconds` has each application join after the measurement starts: a list with a value for each +application, in seconds, where 0 joins at the start, as every application does without the setting. An application +that joins later creates its subscription at the earliest position when the applications start, before the gateways +publish, so that the subscription keeps the backlog until the application opens its pods; it doesn't take part in +the warmup. It has caught up when it first receives a measured message within `applications.caughtUpLatencyMillis` of +its publishing, and the report's Catch-up section shows when each late application joined, its subscription's backlog +then, and how long it took to catch up. The timeout has to leave room for the last application to join and catch up. + `timeoutSeconds` bounds the whole workload: the applications stop waiting for messages after it and the run fails, the gateways give up waiting for a warmup round or the measurement's start, and the launcher waits for the containers a minute longer. It has to be at least the time that the warmup rounds, their delays and the measurement need at the diff --git a/tests/performance/scenarios/iot-telemetry-catch-up.yaml b/tests/performance/scenarios/iot-telemetry-catch-up.yaml new file mode 100644 index 0000000000000..2ad7144ac3d5a --- /dev/null +++ b/tests/performance/scenarios/iot-telemetry-catch-up.yaml @@ -0,0 +1,31 @@ +# 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. + + +# Consumers joining a live stream: the gateways publish at the high-rate scenario's 30,000 msg/s, and 4 of the 5 +# applications join after the measurement starts, 2 of them at the same time. Each late application had its +# subscription from the start, so it catches up on the backlog that built up until it joined, while the gateways keep +# publishing; the report's Catch-up section shows how long each took. The 2 applications that join together read the +# same backlog, so the entries that the first one reads from storage can be served to the other from the broker's +# entry cache. The measurement is 120 s, so the last application joins halfway through it. +extends: iot-telemetry-high-rate.yaml +workloads: + iotTelemetry: + measurement: + messages: 3600000 + applications: + joinSeconds: [0, 20, 20, 40, 60] 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..35005ebb595e3 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 @@ -94,9 +94,28 @@ public record Topics(int count, String prefix) { * @param client the client resources that every pod's client shares * @param env the environment variables of the applications' container, which the launcher sets, such as * {@code PULSAR_MEM} + * @param joinSeconds when each application opens its pods, in seconds after the measurement starts, 0 for at the + * start; empty for every application at the start. An application that joins later has its + * subscription from the start, so it catches up on the backlog that builds up until it joins + * @param caughtUpLatencyMillis an application that joined later has caught up when it first receives a measured + * message within this time of its publishing */ public record Applications(int count, int podsPerApplication, String subscriptionPrefix, Client client, - Map env) { + Map env, List joinSeconds, Integer caughtUpLatencyMillis) { + /** The default of {@code caughtUpLatencyMillis}. */ + public static final int DEFAULT_CAUGHT_UP_LATENCY_MILLIS = 1000; + + public Applications { + joinSeconds = joinSeconds != null ? List.copyOf(joinSeconds) : List.of(); + if (caughtUpLatencyMillis == null) { + caughtUpLatencyMillis = DEFAULT_CAUGHT_UP_LATENCY_MILLIS; + } + } + + public Applications(int count, int podsPerApplication, String subscriptionPrefix, Client client, + Map env) { + this(count, podsPerApplication, subscriptionPrefix, client, env, null, null); + } } /** The applications' client resources: the event loop's and the listener thread pool's threads. */ @@ -176,6 +195,13 @@ 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(applications.joinSeconds().isEmpty() || applications.joinSeconds().size() == applications.count(), + "applications.joinSeconds must have a value for each of the " + applications.count() + + " applications, or none"); + require(applications.joinSeconds().stream().allMatch(seconds -> seconds != null && seconds >= 0 + && seconds < timeoutSeconds), + "applications.joinSeconds must be at least 0 and less than timeoutSeconds"); + require(applications.caughtUpLatencyMillis() >= 1, "applications.caughtUpLatencyMillis 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," @@ -239,6 +265,16 @@ public int podsPerApplication() { return applications.podsPerApplication(); } + /** When the application opens its pods, in seconds after the measurement starts; 0 for at the start. */ + public int joinSeconds(int applicationIndex) { + return applications.joinSeconds().isEmpty() ? 0 : applications.joinSeconds().get(applicationIndex); + } + + /** Whether any application joins after the measurement starts. */ + public boolean hasLateApplications() { + return applications.joinSeconds().stream().anyMatch(seconds -> seconds > 0); + } + public long messageCount() { return Math.addExact(warmupMessageCount(), measurementMessageCount()); } diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java new file mode 100644 index 0000000000000..f66a43cc71a4e --- /dev/null +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java @@ -0,0 +1,62 @@ +/* + * 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 java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.NoSuchFileException; +import java.nio.file.Path; +import java.nio.file.StandardCopyOption; + +/** + * The gateways' measurement start, for the applications that join after it: a file in the coordination directory + * with the start's epoch milliseconds, written atomically. + */ +final class MeasurementStartMarker { + private static final long POLL_INTERVAL_MILLIS = 20; + + private MeasurementStartMarker() { + } + + static void mark(Path directory, String runId, long epochMs) throws IOException { + Files.createDirectories(directory); + Path marker = marker(directory, runId); + Path temporary = marker.resolveSibling(marker.getFileName() + ".tmp"); + Files.writeString(temporary, Long.toString(epochMs), StandardCharsets.UTF_8); + Files.move(temporary, marker, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); + } + + /** Waits for the measurement to start, and returns its epoch milliseconds. */ + static long await(Path directory, String runId, long deadlineNanos) throws Exception { + Path marker = marker(directory, runId); + while (deadlineNanos - System.nanoTime() > 0) { + try { + return Long.parseLong(Files.readString(marker, StandardCharsets.UTF_8).trim()); + } catch (NoSuchFileException e) { + Thread.sleep(POLL_INTERVAL_MILLIS); + } + } + throw new IllegalStateException("Timed out waiting for the measurement to start"); + } + + private static Path marker(Path directory, String runId) { + return directory.resolve("measurement-start-" + runId); + } +} diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java index 3161ef28a7a17..c8f26fe580084 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java @@ -81,6 +81,10 @@ IotScenario scenario() throws Exception { throw new IllegalArgumentException( "Warmup requires the same --run-id for the producer and all consumers"); } + if (scenario.hasLateApplications() && (runId == null || runId.isBlank())) { + throw new IllegalArgumentException( + "Applications that join later require the same --run-id for the producer and all consumers"); + } return scenario; } diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java index 38103fd2e7563..7b72642f098a2 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java @@ -56,6 +56,11 @@ * *

Each application checks its own delivery and ordering, and writes its outputs into a directory named after its * subscription in {@code --output}, as the run report names it. The progress stream sums the applications' counts. + * + *

An application with {@code applications.joinSeconds} after the measurement's start creates its subscription at + * startup, so that the backlog builds up from the start, and opens its pods when it joins. It records when it joined + * and when it caught up: when it first received a measured message within {@code applications.caughtUpLatencyMillis} + * of its publishing. */ @Command(name = "iot-consume", description = "Run the IoT applications, each consuming and validating its own " + "subscription") @@ -89,11 +94,22 @@ public Integer call() throws Exception { sharedResources = SharedClientResources.create(scenario.applications().client().ioThreads(), scenario.applications().client().listenerThreads()); openPods(applications, sharedResources); + for (Application application : applications) { + if (application.joinsLater()) { + // it doesn't take part in the warmup, which the gateways would otherwise wait for + application.skipWarmup(coordinationDirectory(), runId); + } + } phase.set("receiving"); System.out.println("READY applications=" + applications.size() + " clients=" + (long) applications.size() * scenario.podsPerApplication()); + long deadlineNanos = System.nanoTime() + Duration.ofSeconds(scenario.timeoutSeconds()).toNanos(); for (Application application : applications) { - application.startRestarts(sharedResources); + if (application.joinsLater()) { + application.scheduleJoin(sharedResources, coordinationDirectory(), runId, deadlineNanos); + } else { + application.startRestarts(sharedResources); + } } boolean succeeded = receive(scenario, applications); phase.set("finished"); @@ -129,6 +145,7 @@ private boolean receive(IotScenario scenario, List applications) th for (Iterator iterator = receiving.iterator(); iterator.hasNext(); ) { Application application = iterator.next(); application.checkRestarts(); + application.checkJoin(); application.markWarmupRounds(coordinationDirectory(), runId); if (application.receivedEveryMessage()) { succeeded &= application.finish(); @@ -154,7 +171,8 @@ private boolean receive(IotScenario scenario, List applications) th */ private static void openPods(List applications, PulsarClientSharedResources sharedResources) throws Exception { - long pods = (long) applications.size() * applications.get(0).scenario.podsPerApplication(); + long pods = applications.stream().filter(application -> !application.joinsLater()).count() + * applications.get(0).scenario.podsPerApplication(); AtomicLong opened = new AtomicLong(); ExecutorService executor = Executors.newFixedThreadPool(Math.min(applications.size(), MAX_PARALLEL_STARTS), runnable -> { @@ -174,7 +192,11 @@ private static void openPods(List applications, PulsarClientSharedR List> starts = new ArrayList<>(applications.size()); for (Application application : applications) { starts.add(executor.submit(() -> { - application.openPods(sharedResources, opened); + if (application.joinsLater()) { + application.createSubscription(sharedResources); + } else { + application.openPods(sharedResources, opened); + } return null; })); } @@ -232,6 +254,12 @@ private static final class Application { private final HdrLatencyRecorder receiveLatency; private final AtomicLong firstMeasurementReceiptEpochMs = new AtomicLong(); private final AtomicLong lastMeasurementReceiptEpochMs = new AtomicLong(); + // For an application that joins later: when it joined, when it caught up and how many messages it had then + private final AtomicLong joinEpochMs = new AtomicLong(); + private final AtomicLong caughtUpEpochMs = new AtomicLong(); + private final AtomicLong messagesWhenCaughtUp = new AtomicLong(); + private final AtomicReference joinFailure = new AtomicReference<>(); + private Thread joiner; // Guarded by itself, as the restarts replace pods private final List pods; private final AtomicBoolean stopping = new AtomicBoolean(); @@ -254,6 +282,68 @@ HdrLatencyRecorder receiveLatency() { return receiveLatency; } + boolean joinsLater() { + return scenario.joinSeconds(index) > 0; + } + + /** + * Creates the application's subscription on every topic at the earliest position, without keeping a + * consumer, so that it keeps the backlog until the application joins. + */ + void createSubscription(PulsarClientSharedResources sharedResources) throws Exception { + try (PulsarClient client = PulsarClient.builder() + .serviceUrl(scenario.serviceUrl()) + .sharedResources(sharedResources) + .build()) { + client.newConsumer(Schema.BYTES) + .topics(scenario.topicNames()) + .subscriptionName(scenario.subscriptionName(index)) + .subscriptionType(SubscriptionType.Key_Shared) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscribe() + .close(); + } + } + + /** Marks every warmup round received, so that the gateways don't wait for an application that joins later. */ + void skipWarmup(Path coordinationDirectory, String runId) throws IOException { + for (; nextWarmupRound <= scenario.warmupRounds() && scenario.warmupMessageCountPerRound() > 0; + nextWarmupRound++) { + WarmupBarrier.markApplicationComplete(coordinationDirectory, runId, nextWarmupRound, index); + } + } + + /** Opens the application's pods at its join time, counted from the gateways' measurement start. */ + void scheduleJoin(PulsarClientSharedResources sharedResources, Path coordinationDirectory, String runId, + long deadlineNanos) { + joiner = new Thread(() -> { + try { + long measurementStart = MeasurementStartMarker.await(coordinationDirectory, runId, deadlineNanos); + long waitMillis = measurementStart + TimeUnit.SECONDS.toMillis(scenario.joinSeconds(index)) + - System.currentTimeMillis(); + if (waitMillis > 0) { + Thread.sleep(waitMillis); + } + joinEpochMs.set(System.currentTimeMillis()); + System.out.println("JOIN application=" + index + " epochMs=" + joinEpochMs.get()); + openPods(sharedResources, new AtomicLong()); + startRestarts(sharedResources); + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + } catch (Throwable error) { + joinFailure.compareAndSet(null, error); + } + }, "iot-application-join-" + index); + joiner.setDaemon(true); + joiner.start(); + } + + void checkJoin() { + if (joinFailure.get() != null) { + throw new IllegalStateException("Application " + index + " couldn't join", joinFailure.get()); + } + } + void openPods(PulsarClientSharedResources sharedResources, AtomicLong opened) throws Exception { for (int pod = 0; pod < scenario.podsPerApplication(); pod++) { ClientAndConsumer created = createPod(sharedResources, pod); @@ -309,7 +399,11 @@ boolean finish() throws Exception { + " \"firstMeasurementMessageReceivedEpochMs\": " + firstMeasurementReceiptEpochMs.get() + ",\n" + " \"lastMeasurementMessageReceivedEpochMs\": " - + lastMeasurementReceiptEpochMs.get() + "\n}\n"); + + lastMeasurementReceiptEpochMs.get() + ",\n" + + " \"joinSeconds\": " + scenario.joinSeconds(index) + ",\n" + + " \"joinEpochMs\": " + joinEpochMs.get() + ",\n" + + " \"caughtUpEpochMs\": " + caughtUpEpochMs.get() + ",\n" + + " \"messagesWhenCaughtUp\": " + messagesWhenCaughtUp.get() + "\n}\n"); closePods(); finished = true; return summary.valid() && summary.uniqueMessages() == scenario.messageCount(); @@ -318,6 +412,9 @@ boolean finish() throws Exception { /** Closes the application's pods, when it hasn't finished, such as after a failure. */ void close() throws Exception { stopping.set(true); + if (joiner != null) { + joiner.interrupt(); + } if (restarter != null) { restarter.interrupt(); } @@ -326,6 +423,10 @@ void close() throws Exception { private void stopRestarts() throws InterruptedException { stopping.set(true); + if (joiner != null) { + joiner.interrupt(); + joiner.join(TimeUnit.SECONDS.toMillis(10)); + } if (restarter != null) { restarter.interrupt(); restarter.join(TimeUnit.SECONDS.toMillis(10)); @@ -374,6 +475,12 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, (current, received) -> current == 0 ? received : Math.min(current, received)); lastMeasurementReceiptEpochMs.accumulateAndGet(receivedEpochMs, Math::max); + if (joinEpochMs.get() > 0 && caughtUpEpochMs.get() == 0 + && receivedEpochMs - message.getPublishTime() + <= scenario.applications().caughtUpLatencyMillis() + && caughtUpEpochMs.compareAndSet(0, receivedEpochMs)) { + messagesWhenCaughtUp.set(tracker.uniqueMessages()); + } } receiveLatency.recordMillis(receivedEpochMs - message.getPublishTime(), decoded.measurement()); 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..1949015313a15 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 @@ -155,6 +155,10 @@ public Integer call() throws Exception { this.measurementStartEpochMs = measurementStartEpochMs; phase = "measurement"; System.out.println("MEASUREMENT_START epochMs=" + measurementStartEpochMs); + if (scenario.hasLateApplications()) { + // the applications that join later count their join time from it + MeasurementStartMarker.mark(coordinationDirectory(), runId, measurementStartEpochMs); + } } long deviceSequence = deviceSequences[device]++; long producerSequence = producerSequences[producerIndex]++; 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..59085403e953d 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 @@ -97,6 +97,33 @@ public void defaultsMissingWarmupRoundsToOne() { assertThat(scenario.warmupMessageCount()).isEqualTo(1_000); } + @Test + public void joinsApplicationsLater() { + IotScenario scenario = scenario(new IotScenario.Applications(3, 2, "app-", new IotScenario.Client(2, 2), null, + List.of(0, 20, 20), null)); + assertThat(scenario.joinSeconds(0)).isZero(); + assertThat(scenario.joinSeconds(1)).isEqualTo(20); + assertThat(scenario.hasLateApplications()).isTrue(); + assertThat(scenario.applications().caughtUpLatencyMillis()) + .isEqualTo(IotScenario.Applications.DEFAULT_CAUGHT_UP_LATENCY_MILLIS); + // without the setting, every application joins at the start + IotScenario atStart = scenario(new IotScenario.Applications(3, 2, "app-", new IotScenario.Client(2, 2), null)); + assertThat(atStart.joinSeconds(2)).isZero(); + assertThat(atStart.hasLateApplications()).isFalse(); + } + + @Test + public void rejectsInvalidJoinSettings() { + assertThatThrownBy(() -> scenario(new IotScenario.Applications(3, 2, "app-", new IotScenario.Client(2, 2), + null, List.of(0, 20), null))).hasMessageContaining("a value for each of the 3 applications"); + assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), + null, List.of(0, -1), null))).hasMessageContaining("at least 0 and less than timeoutSeconds"); + assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), + null, List.of(0, 300), null))).hasMessageContaining("at least 0 and less than timeoutSeconds"); + assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), + null, List.of(0, 20), 0))).hasMessageContaining("caughtUpLatencyMillis must be at least 1"); + } + @Test public void rejectsTimeBasedWarmupWithoutRateLimit() { assertThatThrownBy(() -> scenario(20, 0, 0, 5_000_000)) @@ -129,6 +156,13 @@ private static IotScenario scenario(int warmupSeconds, long warmupMessages, int 300); } + private static IotScenario scenario(IotScenario.Applications applications) { + return 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), null), + new IotScenario.Topics(2, "persistent://public/default/iot-"), applications, null, 300); + } + private static IotScenario scenario(int warmupSeconds, long warmupMessages, int warmupRounds, int warmupRoundDelaySeconds, int rate, long numberOfMessages, int timeoutSeconds) { diff --git a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java new file mode 100644 index 0000000000000..23e575091139e --- /dev/null +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java @@ -0,0 +1,50 @@ +/* + * 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.nio.file.Files; +import java.nio.file.Path; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import org.testng.annotations.Test; + +public class MeasurementStartMarkerTest { + @Test + public void applicationsAwaitTheGatewaysMeasurementStart() throws Exception { + Path directory = Files.createTempDirectory("coordination"); + CompletableFuture awaited = CompletableFuture.supplyAsync(() -> { + try { + return MeasurementStartMarker.await(directory, "run-1", + System.nanoTime() + TimeUnit.SECONDS.toNanos(10)); + } catch (Exception e) { + throw new IllegalStateException(e); + } + }); + Thread.sleep(50); + assertThat(awaited).isNotDone(); + MeasurementStartMarker.mark(directory, "run-1", 1_790_000_000_000L); + assertThat(awaited.get(10, TimeUnit.SECONDS)).isEqualTo(1_790_000_000_000L); + // another run's marker isn't this run's + assertThatThrownBy(() -> MeasurementStartMarker.await(directory, "run-2", + System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(100))) + .hasMessageContaining("Timed out waiting for the measurement to start"); + } +} From 504e955cedd511a89526f50fc9830a630fb4dd86 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sat, 3 Oct 2026 22:38:32 +0300 Subject: [PATCH 2/5] [improve][test] Report the catch-up's last message and read rate, and join 3 applications 2 s apart The applications that join while the gateways publish at 30,000 msg/s may only catch up after the gateways finish, so the Catch-up section also shows when each received its last message and its read rate since joining. The catch-up scenario joins 3 applications 2 s apart, close enough for the entries that the first reads from storage to stay in the broker's entry cache for the others, and the last one alone at 60 s. Assisted-by: Claude Code (claude-opus-5-5) --- .../tests/performance/report/RunReport.java | 22 +++++++++++-------- .../performance/report/RunReportTest.java | 11 ++++++---- tests/performance/scenarios/README.md | 2 +- .../scenarios/docs/iot-telemetry.md | 13 ++++++----- .../scenarios/iot-telemetry-catch-up.yaml | 13 ++++++----- 5 files changed, 36 insertions(+), 25 deletions(-) diff --git a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java index 666fcd28359da..1e6f11b6550fc 100644 --- a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java +++ b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java @@ -660,15 +660,18 @@ static void appendCatchUp(StringBuilder report, JsonNode workload, List 0 ? String.format(Locale.ROOT, "%.1f s", (caughtUp - joined) / 1000.0) - : "not caught up", - caughtUp > 0 ? String.format(Locale.ROOT, "%,d", messages) : "–", - caughtUp > joined ? String.format(Locale.ROOT, "%,.0f msg/s", - messages * 1000.0 / (caughtUp - joined)) : "–")); + : "not while the gateways published", + lastReceived > joined ? String.format(Locale.ROOT, "%.1f s", (lastReceived - joined) / 1000.0) + : "–", + lastReceived > joined ? String.format(Locale.ROOT, "%,.0f msg/s", + messages * 1000.0 / (lastReceived - joined)) : "–")); } } diff --git a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java index 323d92f9bcc94..8d207b9c46ac4 100644 --- a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java +++ b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java @@ -71,9 +71,11 @@ public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception List consumers = List.of( mapper.readTree("{\"applicationIndex\": 0, \"joinEpochMs\": 0}"), mapper.readTree("{\"applicationIndex\": 1, \"joinEpochMs\": " + (start + 20_000) - + ", \"caughtUpEpochMs\": " + (start + 30_000) + ", \"messagesWhenCaughtUp\": 1500000}"), + + ", \"caughtUpEpochMs\": " + (start + 30_000) + ", \"uniqueMessages\": 4600000," + + " \"lastMeasurementMessageReceivedEpochMs\": " + (start + 120_000) + "}"), mapper.readTree("{\"applicationIndex\": 2, \"joinEpochMs\": " + (start + 40_000) - + ", \"caughtUpEpochMs\": 0, \"messagesWhenCaughtUp\": 0}")); + + ", \"caughtUpEpochMs\": 0, \"uniqueMessages\": 4600000," + + " \"lastMeasurementMessageReceivedEpochMs\": " + (start + 140_000) + "}")); long[] epochs = {start + 19_000, start + 20_500, start + 39_000}; RunReport.Samples samples = new RunReport.Samples(epochs, new double[3], Map.of(), Map.of("app-1", new double[] {570_000, 600_000, 0}, "app-2", new double[] {1_100_000, 1_150_000, @@ -85,8 +87,9 @@ public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception .contains("## Catch-up") .contains("within 500 ms of its publishing") // the backlog of the last sample before the application joined, and its catch-up rate - .contains("| `app-1` | 20.0 s | 570,000 | 10.0 s | 1,500,000 | 150,000 msg/s |") - .contains("| `app-2` | 40.0 s | 1,170,000 | not caught up | – | – |") + .contains("| `app-1` | 20.0 s | 570,000 | 10.0 s | 100.0 s | 46,000 msg/s |") + .contains("| `app-2` | 40.0 s | 1,170,000 | not while the gateways published | 100.0 s" + + " | 46,000 msg/s |") .doesNotContain("`app-0`"); // without late applications, there's no section StringBuilder none = new StringBuilder(); diff --git a/tests/performance/scenarios/README.md b/tests/performance/scenarios/README.md index 0eec53c074510..a62e52d474e17 100644 --- a/tests/performance/scenarios/README.md +++ b/tests/performance/scenarios/README.md @@ -47,7 +47,7 @@ detail. | [`iot-telemetry-small-restarts.yaml`](iot-telemetry-small-restarts.yaml) | The smaller topology, restarting 10 % of each application's clients every 30 s | low, about 3 GB | | [`iot-telemetry-restarts.yaml`](iot-telemetry-restarts.yaml) | The full topology, restarting 10 % of each application's clients every 30 s | medium, about 11 GB | | [`iot-telemetry-high-rate.yaml`](iot-telemetry-high-rate.yaml) | A high rate: five million messages at 30,000 messages per second, from 500 preconnected producers to one topic, consumed by 5 applications with 10 pods each | high, about 14 GB | -| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 120 s, with 4 of its 5 applications joining 20 s (2 of them), 40 s and 60 s after the measurement starts and catching up on their backlogs | high, about 14 GB | +| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 120 s, with 4 of its 5 applications joining 20, 22, 24 and 60 s after the measurement starts and catching up on their backlogs | high, about 14 GB | | [`iot-telemetry-max-rate.yaml`](iot-telemetry-max-rate.yaml) | The maximum rate: the 500 gateways of the high-rate scenario publish four million unbatched messages to one topic without a rate limit, consumed by one application with 20 pods, on single-copy ledgers | high, about 14 GB | ## Configurations diff --git a/tests/performance/scenarios/docs/iot-telemetry.md b/tests/performance/scenarios/docs/iot-telemetry.md index 85046f7584f76..6517f2c9a1859 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -45,11 +45,14 @@ application-visible order across Key_Shared hash-range reassignment. configuration. - [`iot-telemetry-catch-up.yaml`](../iot-telemetry-catch-up.yaml) measures how quickly consumers that join a live stream catch up. The gateways publish at the high-rate scenario's 30,000 messages per second for 120 s, and 4 of - the 5 applications join after the measurement starts: 2 at 20 s, and 1 each at 40 s and 60 s. Each late - application has its subscription from the start, so its backlog builds up until it joins, and it then reads the - backlog while the gateways keep publishing. The report's Catch-up section shows each application's backlog when it - joined and how long it took to catch up. The 2 applications that join together read the same backlog, so the - entries that one of them reads from storage can be served to the other from the broker's entry cache. + the 5 applications join after the measurement starts, 3 of them 2 s apart, at 20, 22 and 24 s, and the last one at + 60 s. Each late application has its subscription from the start, so its backlog builds up until it joins, and it + then reads the backlog while the gateways keep publishing. The report's Catch-up section shows each application's + backlog when it joined and how long it took to catch up. The entries that the first of the 3 reads from storage + are expected to be read by the others, so they stay in the broker's entry cache for its time to live, extended a + few times (`managedLedgerCacheEvictionTimeThresholdMillis` and + `managedLedgerCacheEvictionExtendTTLOfEntriesWithRemainingExpectedReadsMaxTimes`), and the others read them from + there. The last one is too far behind for that, and reads its backlog from storage. - [`iot-telemetry-max-rate.yaml`](../iot-telemetry-max-rate.yaml) runs the high-rate scenario's 500 gateways at the maximum rate: `rate: 0` removes the rate limit, so the gateways publish four million unbatched 128-byte messages as fast as the cluster takes them, with at most 100,000 in flight. One application consumes them with twenty pods diff --git a/tests/performance/scenarios/iot-telemetry-catch-up.yaml b/tests/performance/scenarios/iot-telemetry-catch-up.yaml index 2ad7144ac3d5a..c7285ae752ab6 100644 --- a/tests/performance/scenarios/iot-telemetry-catch-up.yaml +++ b/tests/performance/scenarios/iot-telemetry-catch-up.yaml @@ -17,15 +17,16 @@ # Consumers joining a live stream: the gateways publish at the high-rate scenario's 30,000 msg/s, and 4 of the 5 -# applications join after the measurement starts, 2 of them at the same time. Each late application had its -# subscription from the start, so it catches up on the backlog that built up until it joined, while the gateways keep -# publishing; the report's Catch-up section shows how long each took. The 2 applications that join together read the -# same backlog, so the entries that the first one reads from storage can be served to the other from the broker's -# entry cache. The measurement is 120 s, so the last application joins halfway through it. +# applications join after the measurement starts. Each late application had its subscription from the start, so it +# catches up on the backlog that built up until it joined, while the gateways keep publishing; the report's Catch-up +# section shows how long each took. 3 applications join 2 s apart, at 20, 22 and 24 s, as consumers that scale up or +# restart together do: the entries that the first of them reads from storage stay in the broker's entry cache for the +# others, which are expected to read them within its time to live. The last one joins alone at 60 s, halfway through +# the 120 s measurement, and reads its backlog from storage. extends: iot-telemetry-high-rate.yaml workloads: iotTelemetry: measurement: messages: 3600000 applications: - joinSeconds: [0, 20, 20, 40, 60] + joinSeconds: [0, 20, 22, 24, 60] From e9b948b3685816b6b6abe0947fc8c6ff41eeefcb Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sat, 3 Oct 2026 22:54:36 +0300 Subject: [PATCH 3/5] [improve][test] Address catch-up review: validate join settings, join at measurement start, track catch-up Assisted-by: Claude Code (claude-opus-5-5) --- tests/performance/AGENTS.md | 2 +- tests/performance/docs/run-reports.md | 12 +-- .../launcher/ScenarioFilesTest.java | 8 ++ .../tests/performance/report/RunReport.java | 50 +++++++----- .../performance/report/RunReportTest.java | 49 ++++++++---- tests/performance/scenarios/README.md | 2 +- .../scenarios/docs/iot-telemetry.md | 25 +++--- .../scenarios/iot-telemetry-catch-up.yaml | 7 +- .../performance/tools/CatchUpTracker.java | 79 +++++++++++++++++++ .../tests/performance/tools/IotScenario.java | 12 ++- .../performance/tools/TelemetryConsumer.java | 75 ++++++++++-------- .../performance/tools/CatchUpTrackerTest.java | 54 +++++++++++++ .../performance/tools/IotScenarioTest.java | 30 ++++++- 13 files changed, 314 insertions(+), 91 deletions(-) create mode 100644 tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java create mode 100644 tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java diff --git a/tests/performance/AGENTS.md b/tests/performance/AGENTS.md index 557a4f41ac7d9..b76ce8bfd10da 100644 --- a/tests/performance/AGENTS.md +++ b/tests/performance/AGENTS.md @@ -307,7 +307,7 @@ commit segment can be absent). Use the exact recording names linked from the pro | `console.log.txt` | What the launcher printed, including the progress lines every 10 s | Read | | `launcher.log` | Testcontainers' and the Pulsar containers' logs; kept when the run failed, or with `-Pperformance.keepLauncherLog` | `rg` | | `gateways/gateways-summary.json` | The gateways' counts and throughput; `measurementMessages`, `measurementStartEpochMs` and `measurementEndEpochMs` identify the measured sends | `jq` | -| `applications//application-summary.json` | Each application's unique messages, duplicates, ordering violations and invalid messages; `lastMeasurementMessageReceivedEpochMs` records the last measured receipt | `jq` | +| `applications//application-summary.json` | Each application's unique messages, duplicates, ordering violations and invalid messages; `lastMeasurementMessageReceivedEpochMs` records the last measured receipt; an application that joined after the measurement started (`applications.joinSeconds`) also has `joinEpochMs`, `caughtUpEpochMs`, `caughtUpLatencyMillis` and `messagesWhenCaughtUp` | `jq` | | `applications//ordering-violations.txt` | Samples of the ordering violations; empty in a valid run | Read | | `gateways/container.log.txt`, `applications/container.log.txt` | The workload containers' logs | Read | | `gateways/gateways-latency.hgrm`, `applications//application-latency.hgrm` | The publish and end-to-end latency percentile distributions, in milliseconds, as text | Read | diff --git a/tests/performance/docs/run-reports.md b/tests/performance/docs/run-reports.md index 89366f1edd34a..0dce76e22b766 100644 --- a/tests/performance/docs/run-reports.md +++ b/tests/performance/docs/run-reports.md @@ -118,7 +118,7 @@ didn't receive every message, and the launcher fails a run when an application's ├── heap-dumps/ the heap dumps, when the scenario asks for them, see heap-dumps.md ├── metrics.json the run's metrics in VictoriaMetrics and Grafana, see metrics.md ├── grafana-panels/ panels of Grafana's dashboards over the run, as PNG images -└── coordination/ the warmup barrier markers of the gateways and the applications +└── coordination/ the warmup barrier and measurement start markers of the workloads ``` A profiled run adds the recordings, their flame graphs and a profile report to each profiled component's directory, @@ -144,8 +144,10 @@ report has these sections: - **Throughput**: the gateways' throughput, the delivered throughput until the slowest application received the last message, the measurement's duration and how long the applications were still receiving after the gateways finished. - **Catch-up**, when applications join after the measurement starts (`applications.joinSeconds`): each late - application's join time, its subscription's backlog then, how long it took to catch up, and its catch-up rate. The - backlog chart shows each subscription's backlog building up until its application joins and falling as it catches up. + application's join time, its subscription's backlog then, how long it took to catch up, its catch-up rate and when it + received its last message. An application that didn't catch up shows its rate until its last message instead. When + applications join later, the Throughput section's delivered throughput is mostly the last one's. The backlog chart + shows each subscription's backlog building up until its application joins and falling as it catches up. - **Latency**: the publish latency (send to acknowledgment) and each application's end-to-end latency (publish to consume) at percentiles from p50 to the maximum, with charts by percentile and over time. - **Backlog and rates**: each subscription's backlog and the per-second rates, sampled from the broker's topic @@ -194,12 +196,12 @@ beside it, rendered with [commonmark-java](https://github.com/commonmark/commonm | `gateways/gateways-summary.json` | The gateways' counts and throughput, and the epoch-millisecond boundaries of the measurement | | `gateways/gateways-latency.hdr`, `.hgrm` | The publish latency log, and its percentile distribution in milliseconds, see [Latency logs](#latency-logs) | | `gateways/gateways-state.bin` | The gateways' next sequence number for each device. The launcher compares it with each application's `application-state.bin` and fails the run when they differ, which catches messages missing at the end, where no gap shows | -| `applications//application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt | +| `applications//application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt. An application that joined after the measurement started also has its join time (`joinEpochMs`), when it caught up (`caughtUpEpochMs`, 0 when it didn't) within `caughtUpLatencyMillis`, and its received messages then (`messagesWhenCaughtUp`) | | `applications//application-latency.hdr`, `.hgrm` | The application's end-to-end latency log, and its percentile distribution in milliseconds | | `applications//application-state.bin` | The application's next expected sequence number for each device | | `applications//ordering-violations.txt` | Samples of the ordering violations, with the message ID, topic and receiving thread; empty in a valid run | | `gateways/container.log.txt`, `applications/container.log.txt` | The container's log, named `.txt` so that HTTP servers show it as text | -| `coordination/` | The markers with which the applications tell the gateways that they received a warmup round | +| `coordination/` | The markers with which the applications tell the gateways that they received a warmup round, and the gateways' measurement start, from which the applications that join later count their join time | ## Latency logs diff --git a/tests/performance/launcher/src/test/java/org/apache/pulsar/tests/performance/launcher/ScenarioFilesTest.java b/tests/performance/launcher/src/test/java/org/apache/pulsar/tests/performance/launcher/ScenarioFilesTest.java index f45573e260c96..6163caadd3bff 100644 --- a/tests/performance/launcher/src/test/java/org/apache/pulsar/tests/performance/launcher/ScenarioFilesTest.java +++ b/tests/performance/launcher/src/test/java/org/apache/pulsar/tests/performance/launcher/ScenarioFilesTest.java @@ -102,6 +102,14 @@ private static int ledgerSetting(ClusterSettings cluster, String name) { return value != null ? Integer.parseInt(value) : 2; } + @Test + public void theCatchUpScenarioJoinsApplicationsAfterTheMeasurementStarts() { + ObjectNode applications = (ObjectNode) resolve(SCENARIOS.resolve("iot-telemetry-catch-up.yaml"), List.of()) + .path("workloads").path("iotTelemetry").path("applications"); + assertThat(applications.path("joinSeconds").toString()).isEqualTo("[0,20,22,24,60]"); + assertThat(applications.path("count").asInt()).isEqualTo(5); + } + @Test public void theMemoryConfigurationsSetTheirOwnMemory() { // The configuration that a scenario extends last sets the memory, over the default of iot-telemetry-base.yaml diff --git a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java index 1e6f11b6550fc..f8607eaeb3dd7 100644 --- a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java +++ b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java @@ -324,7 +324,7 @@ public static Path write(Path runDirectory, Run run, ObjectMapper mapper) throws appendThroughput(report, producer, consumers); Path stats = runDirectory.resolve(TOPIC_STATS_FILE); Samples samples = Files.isRegularFile(stats) ? readSamples(stats) : null; - appendCatchUp(report, run.workload(), consumers, samples, measurementStart); + appendCatchUp(report, run.workload(), consumers, samples, measurementStart, measurementEnd); appendLatency(report, runDirectory, run, consumerHistograms, measurementStart); if (samples != null) { appendTopicStats(report, runDirectory, samples, measurementStart, measurementEnd, @@ -642,36 +642,42 @@ private static void appendThroughput(StringBuilder report, JsonNode producer, Li producer.path("measurementElapsedSeconds").asDouble()))) .append(row("Applications still receiving after the gateways finished", String.format(Locale.ROOT, "%.1f s", Math.max(0, lastReceived - end) / 1000.0))); + if (consumers.stream().anyMatch(consumer -> consumer.path("joinEpochMs").asLong() > 0)) { + report.append("\nSome applications joined after the measurement started, so the delivered throughput and" + + " the time the applications were still receiving are mostly those of the last one to catch up;" + + " the Catch-up section has each one's.\n"); + } } /** * The applications that joined after the measurement started ({@code applications.joinSeconds}): when each - * joined, its subscription's backlog then, from the sampled topic stats, and how long it took to catch up. + * joined, its subscription's backlog then, from the sampled topic stats, how long it took to catch up, and its + * catch-up rate. */ static void appendCatchUp(StringBuilder report, JsonNode workload, List consumers, Samples samples, - long measurementStart) { + long measurementStart, long measurementEnd) { List late = consumers.stream().filter(consumer -> consumer.path("joinEpochMs").asLong() > 0) .toList(); if (late.isEmpty()) { return; } - int caughtUpLatency = workload.path("applications").path("caughtUpLatencyMillis").asInt(1000); + int caughtUpLatency = late.get(0).path("caughtUpLatencyMillis") + .asInt(workload.path("applications").path("caughtUpLatencyMillis").asInt(1000)); report.append("\n## Catch-up\n\nThe applications that joined after the measurement started. Each had its" + " subscription from the start, so its backlog built up until it joined. An application has" - + " caught up when it first received a measured message within ") + + " caught up when each of its topics delivered a measured message within ") .append(String.format(Locale.ROOT, "%,d", caughtUpLatency)) - .append(" ms of its publishing, which it may not do before the gateways finish. Its read rate is" - + " every message that it received, warmup included, per second from when it joined until it" - + " received its last message. The backlog is the sampled topic stats' last sample before" - + " the application joined.\n\n| Application | Joined | Backlog when it joined" - + " | Caught up after | Received its last message after | Read rate since joining |\n" - + "|---|---:|---:|---:|---:|---:|\n"); + .append(" ms of its publishing; the time includes opening its pods. Its catch-up rate is the messages" + + " that it received until then, warmup included, per second since it joined; for an" + + " application that didn't catch up, it is the rate until its last message. The backlog is" + + " the sampled topic stats' last sample before the application joined.\n\n" + + "| Application | Joined | Backlog when it joined | Caught up after | Catch-up rate" + + " | Received its last message after |\n|---|---:|---:|---:|---:|---:|\n"); for (JsonNode consumer : late) { String application = applicationName(workload, consumer.path("applicationIndex").asInt()); long joined = consumer.path("joinEpochMs").asLong(); long caughtUp = consumer.path("caughtUpEpochMs").asLong(); long lastReceived = consumer.path("lastMeasurementMessageReceivedEpochMs").asLong(); - long messages = consumer.path("uniqueMessages").asLong(); String backlog = "–"; double[] backlogs = samples != null ? samples.backlog().get(application) : null; if (backlogs != null) { @@ -682,14 +688,22 @@ static void appendCatchUp(StringBuilder report, JsonNode workload, List joined) { + caughtUpAfter = String.format(Locale.ROOT, "%.1f s", (caughtUp - joined) / 1000.0) + + (caughtUp > measurementEnd ? ", after the gateways finished" : ""); + rate = String.format(Locale.ROOT, "%,.0f msg/s", + consumer.path("messagesWhenCaughtUp").asLong() * 1000.0 / (caughtUp - joined)); + } else { + caughtUpAfter = "not caught up"; + rate = lastReceived > joined ? String.format(Locale.ROOT, "%,.0f msg/s until its last message", + consumer.path("uniqueMessages").asLong() * 1000.0 / (lastReceived - joined)) : "–"; + } report.append(String.format(Locale.ROOT, "| `%s` | %.1f s | %s | %s | %s | %s |%n", application, - (joined - measurementStart) / 1000.0, backlog, - caughtUp > 0 ? String.format(Locale.ROOT, "%.1f s", (caughtUp - joined) / 1000.0) - : "not while the gateways published", + (joined - measurementStart) / 1000.0, backlog, caughtUpAfter, rate, lastReceived > joined ? String.format(Locale.ROOT, "%.1f s", (lastReceived - joined) / 1000.0) - : "–", - lastReceived > joined ? String.format(Locale.ROOT, "%,.0f msg/s", - messages * 1000.0 / (lastReceived - joined)) : "–")); + : "–")); } } diff --git a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java index 8d207b9c46ac4..9600712fc4e75 100644 --- a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java +++ b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java @@ -65,38 +65,53 @@ public void deleteRun() throws IOException { @Test public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception { - JsonNode workload = mapper.readTree("{\"applications\": {\"subscriptionPrefix\": \"app-\"," - + " \"caughtUpLatencyMillis\": 500}}"); + JsonNode workload = mapper.readTree("{\"applications\": {\"subscriptionPrefix\": \"app-\"}}"); long start = 1_000_000; + long end = start + 120_000; List consumers = List.of( mapper.readTree("{\"applicationIndex\": 0, \"joinEpochMs\": 0}"), - mapper.readTree("{\"applicationIndex\": 1, \"joinEpochMs\": " + (start + 20_000) - + ", \"caughtUpEpochMs\": " + (start + 30_000) + ", \"uniqueMessages\": 4600000," - + " \"lastMeasurementMessageReceivedEpochMs\": " + (start + 120_000) + "}"), - mapper.readTree("{\"applicationIndex\": 2, \"joinEpochMs\": " + (start + 40_000) - + ", \"caughtUpEpochMs\": 0, \"uniqueMessages\": 4600000," - + " \"lastMeasurementMessageReceivedEpochMs\": " + (start + 140_000) + "}")); - long[] epochs = {start + 19_000, start + 20_500, start + 39_000}; + // caught up while the gateways published + late(1, start + 20_000, start + 30_000, 1_500_000, start + 120_000), + // caught up after the gateways finished + late(2, start + 40_000, start + 130_000, 4_000_000, start + 140_000), + // never caught up, with no sample before it joined + late(3, start + 500, 0, 4_600_000, start + 150_000), + // joined but never received a measured message + late(4, start + 60_000, 0, 0, 0)); + long[] epochs = {start + 1_000, start + 19_000, start + 39_000}; RunReport.Samples samples = new RunReport.Samples(epochs, new double[3], Map.of(), - Map.of("app-1", new double[] {570_000, 600_000, 0}, "app-2", new double[] {1_100_000, 1_150_000, - 1_170_000})); + Map.of("app-1", new double[] {0, 570_000, 0}, "app-2", new double[] {0, 0, 1_170_000}, + "app-3", new double[] {0, 0, 0})); StringBuilder report = new StringBuilder(); - RunReport.appendCatchUp(report, workload, consumers, samples, start); + RunReport.appendCatchUp(report, workload, consumers, samples, start, end); assertThat(report.toString()) .contains("## Catch-up") + // the threshold that the application wrote .contains("within 500 ms of its publishing") - // the backlog of the last sample before the application joined, and its catch-up rate - .contains("| `app-1` | 20.0 s | 570,000 | 10.0 s | 100.0 s | 46,000 msg/s |") - .contains("| `app-2` | 40.0 s | 1,170,000 | not while the gateways published | 100.0 s" - + " | 46,000 msg/s |") + .contains("| `app-1` | 20.0 s | 570,000 | 10.0 s | 150,000 msg/s | 100.0 s |") + .contains("| `app-2` | 40.0 s | 1,170,000 | 90.0 s, after the gateways finished | 44,444 msg/s" + + " | 100.0 s |") + .contains("| `app-3` | 0.5 s | – | not caught up | 30,769 msg/s until its last message | 149.5 s |") + .contains("| `app-4` | 60.0 s | – | not caught up | – | – |") .doesNotContain("`app-0`"); + // without topic stats + StringBuilder withoutSamples = new StringBuilder(); + RunReport.appendCatchUp(withoutSamples, workload, consumers.subList(0, 2), null, start, end); + assertThat(withoutSamples.toString()).contains("| `app-1` | 20.0 s | – | 10.0 s |"); // without late applications, there's no section StringBuilder none = new StringBuilder(); - RunReport.appendCatchUp(none, workload, consumers.subList(0, 1), samples, start); + RunReport.appendCatchUp(none, workload, consumers.subList(0, 1), samples, start, end); assertThat(none.toString()).isEmpty(); } + private JsonNode late(int index, long joined, long caughtUp, long messages, long lastReceived) throws Exception { + return mapper.readTree("{\"applicationIndex\": " + index + ", \"joinEpochMs\": " + joined + + ", \"caughtUpLatencyMillis\": 500, \"caughtUpEpochMs\": " + caughtUp + + ", \"messagesWhenCaughtUp\": " + messages + ", \"uniqueMessages\": 4600000" + + ", \"lastMeasurementMessageReceivedEpochMs\": " + lastReceived + "}"); + } + @Test public void reportsCorrectnessThroughputLatencyAndSampledStats() throws IOException { Files.createDirectories(run.resolve("gateways")); diff --git a/tests/performance/scenarios/README.md b/tests/performance/scenarios/README.md index a62e52d474e17..e3fe215ed5c52 100644 --- a/tests/performance/scenarios/README.md +++ b/tests/performance/scenarios/README.md @@ -47,7 +47,7 @@ detail. | [`iot-telemetry-small-restarts.yaml`](iot-telemetry-small-restarts.yaml) | The smaller topology, restarting 10 % of each application's clients every 30 s | low, about 3 GB | | [`iot-telemetry-restarts.yaml`](iot-telemetry-restarts.yaml) | The full topology, restarting 10 % of each application's clients every 30 s | medium, about 11 GB | | [`iot-telemetry-high-rate.yaml`](iot-telemetry-high-rate.yaml) | A high rate: five million messages at 30,000 messages per second, from 500 preconnected producers to one topic, consumed by 5 applications with 10 pods each | high, about 14 GB | -| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 120 s, with 4 of its 5 applications joining 20, 22, 24 and 60 s after the measurement starts and catching up on their backlogs | high, about 14 GB | +| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 180 s, with 4 of its 5 applications joining 20, 22, 24 and 60 s after the measurement starts and catching up on their backlogs | high, about 14 GB | | [`iot-telemetry-max-rate.yaml`](iot-telemetry-max-rate.yaml) | The maximum rate: the 500 gateways of the high-rate scenario publish four million unbatched messages to one topic without a rate limit, consumed by one application with 20 pods, on single-copy ledgers | high, about 14 GB | ## Configurations diff --git a/tests/performance/scenarios/docs/iot-telemetry.md b/tests/performance/scenarios/docs/iot-telemetry.md index 6517f2c9a1859..8b507044433cc 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -44,15 +44,20 @@ application-visible order across Key_Shared hash-range reassignment. applications receive, and an application that falls behind can stall for tens of seconds. It uses the high-memory configuration. - [`iot-telemetry-catch-up.yaml`](../iot-telemetry-catch-up.yaml) measures how quickly consumers that join a live - stream catch up. The gateways publish at the high-rate scenario's 30,000 messages per second for 120 s, and 4 of + stream catch up. The gateways publish at the high-rate scenario's 30,000 messages per second for 180 s, and 4 of the 5 applications join after the measurement starts, 3 of them 2 s apart, at 20, 22 and 24 s, and the last one at 60 s. Each late application has its subscription from the start, so its backlog builds up until it joins, and it then reads the backlog while the gateways keep publishing. The report's Catch-up section shows each application's - backlog when it joined and how long it took to catch up. The entries that the first of the 3 reads from storage - are expected to be read by the others, so they stay in the broker's entry cache for its time to live, extended a - few times (`managedLedgerCacheEvictionTimeThresholdMillis` and - `managedLedgerCacheEvictionExtendTTLOfEntriesWithRemainingExpectedReadsMaxTimes`), and the others read them from - there. The last one is too far behind for that, and reads its backlog from storage. + backlog when it joined and how long it took to catch up; the last one may only catch up after the gateways finish. + - **The broker's entry cache:** a subscription without consumers isn't an active cursor, so nothing keeps the + entries for a late application before it joins. Once the 3 that join 2 s apart are reading, the entries that the + first of them reads from storage are expected to be read by the others behind it, so they stay in the cache for + its time to live, `managedLedgerCacheEvictionTimeThresholdMillis` (1 s), extended up to + `managedLedgerCacheEvictionExtendTTLOfEntriesWithRemainingExpectedReadsMaxTimes` (5) times, about 6 s, which the + 2 s spacing is within. The others then read them from the cache. The last application is too far behind for that, + and reads its backlog from storage. + - **Where reads come from:** the managed ledger's cache hit and miss rates (`pulsar_ml_cache_hits_rate`, + `pulsar_ml_cache_misses_rate`) in the metrics show which reads the cache served. - [`iot-telemetry-max-rate.yaml`](../iot-telemetry-max-rate.yaml) runs the high-rate scenario's 500 gateways at the maximum rate: `rate: 0` removes the rate limit, so the gateways publish four million unbatched 128-byte messages as fast as the cluster takes them, with at most 100,000 in flight. One application consumes them with twenty pods @@ -184,9 +189,11 @@ workloads: application, in seconds, where 0 joins at the start, as every application does without the setting. An application that joins later creates its subscription at the earliest position when the applications start, before the gateways publish, so that the subscription keeps the backlog until the application opens its pods; it doesn't take part in -the warmup. It has caught up when it first receives a measured message within `applications.caughtUpLatencyMillis` of -its publishing, and the report's Catch-up section shows when each late application joined, its subscription's backlog -then, and how long it took to catch up. The timeout has to leave room for the last application to join and catch up. +the warmup. It has caught up when each of its topics has delivered a measured message within +`applications.caughtUpLatencyMillis` of its publishing; the time includes opening its pods, and with Key_Shared +redeliveries a few older messages may still arrive after it. The report's Catch-up section shows when each late +application joined, its subscription's backlog then, and how long it took to catch up. The warmup and the latest join +have to be within the timeout, which also has to leave room for the last application to catch up. `timeoutSeconds` bounds the whole workload: the applications stop waiting for messages after it and the run fails, the gateways give up waiting for a warmup round or the measurement's start, and the launcher waits for the containers diff --git a/tests/performance/scenarios/iot-telemetry-catch-up.yaml b/tests/performance/scenarios/iot-telemetry-catch-up.yaml index c7285ae752ab6..d68eeabed9d23 100644 --- a/tests/performance/scenarios/iot-telemetry-catch-up.yaml +++ b/tests/performance/scenarios/iot-telemetry-catch-up.yaml @@ -21,12 +21,13 @@ # catches up on the backlog that built up until it joined, while the gateways keep publishing; the report's Catch-up # section shows how long each took. 3 applications join 2 s apart, at 20, 22 and 24 s, as consumers that scale up or # restart together do: the entries that the first of them reads from storage stay in the broker's entry cache for the -# others, which are expected to read them within its time to live. The last one joins alone at 60 s, halfway through -# the 120 s measurement, and reads its backlog from storage. +# others, which are expected to read them within its time to live. The last one joins alone at 60 s and reads its +# backlog from storage. The measurement is 180 s, so that the applications that join at 20 to 24 s catch up while the +# gateways publish. extends: iot-telemetry-high-rate.yaml workloads: iotTelemetry: measurement: - messages: 3600000 + messages: 5400000 applications: joinSeconds: [0, 20, 22, 24, 60] diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java new file mode 100644 index 0000000000000..9e84d0e8a37f8 --- /dev/null +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java @@ -0,0 +1,79 @@ +/* + * 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 java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; + +/** + * When an application that joined after the measurement started catches up: when each of its topics has delivered a + * measured message received within the threshold of its publishing, after the application joined. An application + * that joined at the start doesn't catch up. + */ +final class CatchUpTracker { + private final int topicCount; + private final long thresholdMillis; + private final Set caughtUpTopics = ConcurrentHashMap.newKeySet(); + private final AtomicLong joinEpochMs = new AtomicLong(); + private final AtomicLong caughtUpEpochMs = new AtomicLong(); + private final AtomicLong messagesWhenCaughtUp = new AtomicLong(); + + CatchUpTracker(int topicCount, long thresholdMillis) { + this.topicCount = topicCount; + this.thresholdMillis = thresholdMillis; + } + + void joined(long epochMs) { + joinEpochMs.set(epochMs); + } + + /** + * Records a received measured message. + * + * @param messages the application's received messages, read when it has caught up + */ + void received(String topic, long publishEpochMs, long receivedEpochMs, LongSupplier messages) { + if (joinEpochMs.get() == 0 || caughtUpEpochMs.get() != 0 + || receivedEpochMs - publishEpochMs > thresholdMillis) { + return; + } + if (caughtUpTopics.add(topic) && caughtUpTopics.size() >= topicCount + && caughtUpEpochMs.compareAndSet(0, receivedEpochMs)) { + messagesWhenCaughtUp.set(messages.getAsLong()); + } + } + + long joinEpochMs() { + return joinEpochMs.get(); + } + + long caughtUpEpochMs() { + return caughtUpEpochMs.get(); + } + + long messagesWhenCaughtUp() { + return messagesWhenCaughtUp.get(); + } + + long thresholdMillis() { + return thresholdMillis; + } +} 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 35005ebb595e3..2559b7c9460bf 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 @@ -19,6 +19,7 @@ package org.apache.pulsar.tests.performance.tools; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Locale; import java.util.Map; @@ -106,7 +107,8 @@ public record Applications(int count, int podsPerApplication, String subscriptio public static final int DEFAULT_CAUGHT_UP_LATENCY_MILLIS = 1000; public Applications { - joinSeconds = joinSeconds != null ? List.copyOf(joinSeconds) : List.of(); + // kept with any nulls, which the scenario's validation rejects with its message + joinSeconds = joinSeconds != null ? Collections.unmodifiableList(new ArrayList<>(joinSeconds)) : List.of(); if (caughtUpLatencyMillis == null) { caughtUpLatencyMillis = DEFAULT_CAUGHT_UP_LATENCY_MILLIS; } @@ -198,12 +200,14 @@ boolean enabled() { require(applications.joinSeconds().isEmpty() || applications.joinSeconds().size() == applications.count(), "applications.joinSeconds must have a value for each of the " + applications.count() + " applications, or none"); + // the joins count from the measurement's start, after the warmup + long warmupTotalSeconds = minimumRuntimeSeconds - durationSeconds; require(applications.joinSeconds().stream().allMatch(seconds -> seconds != null && seconds >= 0 - && seconds < timeoutSeconds), - "applications.joinSeconds must be at least 0 and less than timeoutSeconds"); + && warmupTotalSeconds + seconds < timeoutSeconds), + "applications.joinSeconds must be at least 0, and the warmup (" + warmupTotalSeconds + + " s) and the latest join must be within timeoutSeconds"); require(applications.caughtUpLatencyMillis() >= 1, "applications.caughtUpLatencyMillis 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," + " but the workload needs %d s: %d s of warmup (%d round(s) of %d s%s) and %d s of" + " measurement%s. The applications stop waiting at the timeout, so set timeoutSeconds to" diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java index 7b72642f098a2..4062768834a6b 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java @@ -101,8 +101,10 @@ public Integer call() throws Exception { } } phase.set("receiving"); + // The clients open now; the late applications' pods open when they join System.out.println("READY applications=" + applications.size() + " clients=" - + (long) applications.size() * scenario.podsPerApplication()); + + applications.stream().filter(application -> !application.joinsLater()).count() + * scenario.podsPerApplication()); long deadlineNanos = System.nanoTime() + Duration.ofSeconds(scenario.timeoutSeconds()).toNanos(); for (Application application : applications) { if (application.joinsLater()) { @@ -159,6 +161,7 @@ private boolean receive(IotScenario scenario, List applications) th // The applications that timed out write what they received too for (Application application : receiving) { application.checkRestarts(); + application.checkJoin(); succeeded &= application.finish(); } return succeeded; @@ -254,17 +257,15 @@ private static final class Application { private final HdrLatencyRecorder receiveLatency; private final AtomicLong firstMeasurementReceiptEpochMs = new AtomicLong(); private final AtomicLong lastMeasurementReceiptEpochMs = new AtomicLong(); - // For an application that joins later: when it joined, when it caught up and how many messages it had then - private final AtomicLong joinEpochMs = new AtomicLong(); - private final AtomicLong caughtUpEpochMs = new AtomicLong(); - private final AtomicLong messagesWhenCaughtUp = new AtomicLong(); + // For an application that joins later: when it joined and when it caught up + private final CatchUpTracker catchUp; private final AtomicReference joinFailure = new AtomicReference<>(); - private Thread joiner; + private volatile Thread joiner; // Guarded by itself, as the restarts replace pods private final List pods; private final AtomicBoolean stopping = new AtomicBoolean(); private final AtomicReference restartFailure = new AtomicReference<>(); - private Thread restarter; + private volatile Thread restarter; private int nextWarmupRound = 1; private volatile boolean finished; @@ -273,6 +274,7 @@ private static final class Application { this.index = index; this.output = output; tracker = new DeviceSequenceTracker(scenario.deviceCount()); + catchUp = new CatchUpTracker(scenario.topicCount(), scenario.applications().caughtUpLatencyMillis()); pods = new ArrayList<>(scenario.podsPerApplication()); receiveLatency = new HdrLatencyRecorder(output.resolve("application-latency.hdr"), PerformanceTool.MAX_LATENCY_MICROS); @@ -324,8 +326,10 @@ void scheduleJoin(PulsarClientSharedResources sharedResources, Path coordination if (waitMillis > 0) { Thread.sleep(waitMillis); } - joinEpochMs.set(System.currentTimeMillis()); - System.out.println("JOIN application=" + index + " epochMs=" + joinEpochMs.get()); + // when it starts to open its pods, so the catch-up includes connecting them + long joinEpochMs = System.currentTimeMillis(); + catchUp.joined(joinEpochMs); + System.out.println("JOIN application=" + index + " epochMs=" + joinEpochMs); openPods(sharedResources, new AtomicLong()); startRestarts(sharedResources); } catch (InterruptedException interrupted) { @@ -348,6 +352,11 @@ void openPods(PulsarClientSharedResources sharedResources, AtomicLong opened) th for (int pod = 0; pod < scenario.podsPerApplication(); pod++) { ClientAndConsumer created = createPod(sharedResources, pod); synchronized (pods) { + // a late application's joiner can open a pod after the application stopped + if (stopping.get()) { + created.close(); + return; + } pods.add(created); } opened.incrementAndGet(); @@ -355,9 +364,10 @@ void openPods(PulsarClientSharedResources sharedResources, AtomicLong opened) th } void startRestarts(PulsarClientSharedResources sharedResources) { - if (scenario.behaviors().podRestarts().enabled()) { - restarter = new Thread(() -> restartPods(sharedResources), "iot-client-restarter-" + index); - restarter.start(); + if (scenario.behaviors().podRestarts().enabled() && !stopping.get()) { + Thread thread = new Thread(() -> restartPods(sharedResources), "iot-client-restarter-" + index); + restarter = thread; + thread.start(); } } @@ -401,9 +411,10 @@ boolean finish() throws Exception { + " \"lastMeasurementMessageReceivedEpochMs\": " + lastMeasurementReceiptEpochMs.get() + ",\n" + " \"joinSeconds\": " + scenario.joinSeconds(index) + ",\n" - + " \"joinEpochMs\": " + joinEpochMs.get() + ",\n" - + " \"caughtUpEpochMs\": " + caughtUpEpochMs.get() + ",\n" - + " \"messagesWhenCaughtUp\": " + messagesWhenCaughtUp.get() + "\n}\n"); + + " \"joinEpochMs\": " + catchUp.joinEpochMs() + ",\n" + + " \"caughtUpLatencyMillis\": " + catchUp.thresholdMillis() + ",\n" + + " \"caughtUpEpochMs\": " + catchUp.caughtUpEpochMs() + ",\n" + + " \"messagesWhenCaughtUp\": " + catchUp.messagesWhenCaughtUp() + "\n}\n"); closePods(); finished = true; return summary.valid() && summary.uniqueMessages() == scenario.messageCount(); @@ -412,24 +423,30 @@ boolean finish() throws Exception { /** Closes the application's pods, when it hasn't finished, such as after a failure. */ void close() throws Exception { stopping.set(true); - if (joiner != null) { - joiner.interrupt(); + Thread joining = joiner; + if (joining != null) { + joining.interrupt(); + // so that it doesn't open pods after they're closed, or use the shared resources after they are + joining.join(TimeUnit.SECONDS.toMillis(10)); } - if (restarter != null) { - restarter.interrupt(); + Thread restarting = restarter; + if (restarting != null) { + restarting.interrupt(); } closePods(); } private void stopRestarts() throws InterruptedException { stopping.set(true); - if (joiner != null) { - joiner.interrupt(); - joiner.join(TimeUnit.SECONDS.toMillis(10)); + Thread joining = joiner; + if (joining != null) { + joining.interrupt(); + joining.join(TimeUnit.SECONDS.toMillis(10)); } - if (restarter != null) { - restarter.interrupt(); - restarter.join(TimeUnit.SECONDS.toMillis(10)); + Thread restarting = restarter; + if (restarting != null) { + restarting.interrupt(); + restarting.join(TimeUnit.SECONDS.toMillis(10)); } } @@ -475,12 +492,8 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, (current, received) -> current == 0 ? received : Math.min(current, received)); lastMeasurementReceiptEpochMs.accumulateAndGet(receivedEpochMs, Math::max); - if (joinEpochMs.get() > 0 && caughtUpEpochMs.get() == 0 - && receivedEpochMs - message.getPublishTime() - <= scenario.applications().caughtUpLatencyMillis() - && caughtUpEpochMs.compareAndSet(0, receivedEpochMs)) { - messagesWhenCaughtUp.set(tracker.uniqueMessages()); - } + catchUp.received(message.getTopicName(), message.getPublishTime(), + receivedEpochMs, tracker::uniqueMessages); } receiveLatency.recordMillis(receivedEpochMs - message.getPublishTime(), decoded.measurement()); diff --git a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java new file mode 100644 index 0000000000000..c95c860454ea4 --- /dev/null +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java @@ -0,0 +1,54 @@ +/* + * 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 org.testng.annotations.Test; + +public class CatchUpTrackerTest { + @Test + public void catchesUpWhenEveryTopicDeliveredWithinTheThresholdAfterJoining() { + CatchUpTracker tracker = new CatchUpTracker(2, 1_000); + // before the application joins, nothing counts + tracker.received("t0", 10_000, 10_100, () -> 1); + tracker.joined(20_000); + assertThat(tracker.caughtUpEpochMs()).isZero(); + // an old message, beyond the threshold, then one topic caught up at the threshold's boundary + tracker.received("t0", 1_000, 21_000, () -> 2); + tracker.received("t0", 21_000, 22_000, () -> 3); + assertThat(tracker.caughtUpEpochMs()).isZero(); + // the other topic too: caught up, with the messages received by then + tracker.received("t1", 22_500, 23_000, () -> 400); + assertThat(tracker.caughtUpEpochMs()).isEqualTo(23_000); + assertThat(tracker.messagesWhenCaughtUp()).isEqualTo(400); + // the first catch-up counts + tracker.received("t1", 30_000, 30_001, () -> 500); + assertThat(tracker.caughtUpEpochMs()).isEqualTo(23_000); + assertThat(tracker.messagesWhenCaughtUp()).isEqualTo(400); + assertThat(tracker.joinEpochMs()).isEqualTo(20_000); + assertThat(tracker.thresholdMillis()).isEqualTo(1_000); + } + + @Test + public void anApplicationThatJoinedAtTheStartDoesntCatchUp() { + CatchUpTracker tracker = new CatchUpTracker(1, 1_000); + tracker.received("t0", 1_000, 1_010, () -> 1); + assertThat(tracker.caughtUpEpochMs()).isZero(); + } +} 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 59085403e953d..ad1ce1dfedeec 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 @@ -23,6 +23,7 @@ import java.nio.file.Files; import java.nio.file.Path; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import java.util.Map; import org.apache.pulsar.tests.performance.common.YamlScenarioLoader; @@ -97,6 +98,29 @@ public void defaultsMissingWarmupRoundsToOne() { assertThat(scenario.warmupMessageCount()).isEqualTo(1_000); } + @Test + public void applicationsThatJoinLaterRequireARunId() throws Exception { + Path config = Files.createTempFile("iot-scenario", ".yaml"); + try { + YamlScenarioLoader loader = new YamlScenarioLoader(); + loader.mapper().writeValue(config.toFile(), Map.of("workloads", Map.of("iotTelemetry", + scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), null, + List.of(0, 20), null))))); + List arguments = new ArrayList<>(List.of("iot-consume", "--config", config.toString(), + "--output", "results")); + var parsed = new CommandLine(new PerformanceTool()).parseArgs(arguments.toArray(String[]::new)); + var tool = (PerformanceTool.ScenarioCommand) parsed.subcommand().commandSpec().userObject(); + assertThatThrownBy(tool::scenario).isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("join later require the same --run-id"); + arguments.addAll(List.of("--run-id", "shared-run")); + parsed = new CommandLine(new PerformanceTool()).parseArgs(arguments.toArray(String[]::new)); + tool = (PerformanceTool.ScenarioCommand) parsed.subcommand().commandSpec().userObject(); + assertThat(tool.scenario().joinSeconds(1)).isEqualTo(20); + } finally { + Files.deleteIfExists(config); + } + } + @Test public void joinsApplicationsLater() { IotScenario scenario = scenario(new IotScenario.Applications(3, 2, "app-", new IotScenario.Client(2, 2), null, @@ -117,9 +141,11 @@ public void rejectsInvalidJoinSettings() { assertThatThrownBy(() -> scenario(new IotScenario.Applications(3, 2, "app-", new IotScenario.Client(2, 2), null, List.of(0, 20), null))).hasMessageContaining("a value for each of the 3 applications"); assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), - null, List.of(0, -1), null))).hasMessageContaining("at least 0 and less than timeoutSeconds"); + null, List.of(0, -1), null))).hasMessageContaining("must be at least 0"); + assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), + null, List.of(0, 300), null))).hasMessageContaining("the latest join must be within timeoutSeconds"); assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), - null, List.of(0, 300), null))).hasMessageContaining("at least 0 and less than timeoutSeconds"); + null, Arrays.asList(0, null), null))).hasMessageContaining("applications.joinSeconds must be"); assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), null, List.of(0, 20), 0))).hasMessageContaining("caughtUpLatencyMillis must be at least 1"); } From 4c2f347b63565603afbeaba54d6b756572590f89 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sun, 4 Oct 2026 01:07:13 +0300 Subject: [PATCH 4/5] [improve][test] Address the follow-up catch-up review: latest topic receipt, pod cleanup on shutdown, report edge cases Assisted-by: Claude Code (claude-opus-5-5) --- tests/performance/AGENTS.md | 2 +- tests/performance/docs/run-reports.md | 2 +- .../tests/performance/report/RunReport.java | 8 +-- .../performance/report/RunReportTest.java | 11 ++-- .../scenarios/docs/iot-telemetry.md | 3 +- .../performance/tools/CatchUpTracker.java | 15 ++++-- .../tests/performance/tools/IotScenario.java | 4 +- .../performance/tools/PerformanceTool.java | 3 +- .../performance/tools/TelemetryConsumer.java | 53 ++++++++++--------- .../performance/tools/CatchUpTrackerTest.java | 10 ++++ .../performance/tools/IotScenarioTest.java | 16 +++++- 11 files changed, 82 insertions(+), 45 deletions(-) diff --git a/tests/performance/AGENTS.md b/tests/performance/AGENTS.md index b76ce8bfd10da..6c3b6345dfb55 100644 --- a/tests/performance/AGENTS.md +++ b/tests/performance/AGENTS.md @@ -307,7 +307,7 @@ commit segment can be absent). Use the exact recording names linked from the pro | `console.log.txt` | What the launcher printed, including the progress lines every 10 s | Read | | `launcher.log` | Testcontainers' and the Pulsar containers' logs; kept when the run failed, or with `-Pperformance.keepLauncherLog` | `rg` | | `gateways/gateways-summary.json` | The gateways' counts and throughput; `measurementMessages`, `measurementStartEpochMs` and `measurementEndEpochMs` identify the measured sends | `jq` | -| `applications//application-summary.json` | Each application's unique messages, duplicates, ordering violations and invalid messages; `lastMeasurementMessageReceivedEpochMs` records the last measured receipt; an application that joined after the measurement started (`applications.joinSeconds`) also has `joinEpochMs`, `caughtUpEpochMs`, `caughtUpLatencyMillis` and `messagesWhenCaughtUp` | `jq` | +| `applications//application-summary.json` | Each application's unique messages, duplicates, ordering violations and invalid messages; `lastMeasurementMessageReceivedEpochMs` records the last measured receipt; `joinSeconds`, `joinEpochMs`, `caughtUpEpochMs`, `caughtUpLatencyMillis` and `messagesWhenCaughtUp` describe an application that joined after the measurement started (`applications.joinSeconds`), with 0 times otherwise | `jq` | | `applications//ordering-violations.txt` | Samples of the ordering violations; empty in a valid run | Read | | `gateways/container.log.txt`, `applications/container.log.txt` | The workload containers' logs | Read | | `gateways/gateways-latency.hgrm`, `applications//application-latency.hgrm` | The publish and end-to-end latency percentile distributions, in milliseconds, as text | Read | diff --git a/tests/performance/docs/run-reports.md b/tests/performance/docs/run-reports.md index 0dce76e22b766..9b7927cc87092 100644 --- a/tests/performance/docs/run-reports.md +++ b/tests/performance/docs/run-reports.md @@ -196,7 +196,7 @@ beside it, rendered with [commonmark-java](https://github.com/commonmark/commonm | `gateways/gateways-summary.json` | The gateways' counts and throughput, and the epoch-millisecond boundaries of the measurement | | `gateways/gateways-latency.hdr`, `.hgrm` | The publish latency log, and its percentile distribution in milliseconds, see [Latency logs](#latency-logs) | | `gateways/gateways-state.bin` | The gateways' next sequence number for each device. The launcher compares it with each application's `application-state.bin` and fails the run when they differ, which catches messages missing at the end, where no gap shows | -| `applications//application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt. An application that joined after the measurement started also has its join time (`joinEpochMs`), when it caught up (`caughtUpEpochMs`, 0 when it didn't) within `caughtUpLatencyMillis`, and its received messages then (`messagesWhenCaughtUp`) | +| `applications//application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt. Its join setting (`joinSeconds`), and for an application that joined after the measurement started, its join time (`joinEpochMs`), when it caught up (`caughtUpEpochMs`) within `caughtUpLatencyMillis`, and its received messages then (`messagesWhenCaughtUp`); the times are 0 for an application that joined at the start or didn't catch up | | `applications//application-latency.hdr`, `.hgrm` | The application's end-to-end latency log, and its percentile distribution in milliseconds | | `applications//application-state.bin` | The application's next expected sequence number for each device | | `applications//ordering-violations.txt` | Samples of the ordering violations, with the message ID, topic and receiving thread; empty in a valid run | diff --git a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java index f8607eaeb3dd7..012c8db68b7f1 100644 --- a/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java +++ b/tests/performance/report-tool/src/main/java/org/apache/pulsar/tests/performance/report/RunReport.java @@ -644,7 +644,7 @@ private static void appendThroughput(StringBuilder report, JsonNode producer, Li String.format(Locale.ROOT, "%.1f s", Math.max(0, lastReceived - end) / 1000.0))); if (consumers.stream().anyMatch(consumer -> consumer.path("joinEpochMs").asLong() > 0)) { report.append("\nSome applications joined after the measurement started, so the delivered throughput and" - + " the time the applications were still receiving are mostly those of the last one to catch up;" + + " the time the applications were still receiving are mostly those of the last one to finish;" + " the Catch-up section has each one's.\n"); } } @@ -690,11 +690,11 @@ static void appendCatchUp(StringBuilder report, JsonNode workload, List joined) { + if (caughtUp > 0 && caughtUp >= joined) { caughtUpAfter = String.format(Locale.ROOT, "%.1f s", (caughtUp - joined) / 1000.0) + (caughtUp > measurementEnd ? ", after the gateways finished" : ""); - rate = String.format(Locale.ROOT, "%,.0f msg/s", - consumer.path("messagesWhenCaughtUp").asLong() * 1000.0 / (caughtUp - joined)); + rate = caughtUp > joined ? String.format(Locale.ROOT, "%,.0f msg/s", + consumer.path("messagesWhenCaughtUp").asLong() * 1000.0 / (caughtUp - joined)) : "–"; } else { caughtUpAfter = "not caught up"; rate = lastReceived > joined ? String.format(Locale.ROOT, "%,.0f msg/s until its last message", diff --git a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java index 9600712fc4e75..53ea38d01bc21 100644 --- a/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java +++ b/tests/performance/report-tool/src/test/java/org/apache/pulsar/tests/performance/report/RunReportTest.java @@ -72,12 +72,14 @@ public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception mapper.readTree("{\"applicationIndex\": 0, \"joinEpochMs\": 0}"), // caught up while the gateways published late(1, start + 20_000, start + 30_000, 1_500_000, start + 120_000), - // caught up after the gateways finished - late(2, start + 40_000, start + 130_000, 4_000_000, start + 140_000), + // caught up with the gateways' last messages, within the threshold after they finished + late(2, start + 40_000, start + 120_400, 4_000_000, start + 140_000), // never caught up, with no sample before it joined late(3, start + 500, 0, 4_600_000, start + 150_000), // joined but never received a measured message - late(4, start + 60_000, 0, 0, 0)); + late(4, start + 60_000, 0, 0, 0), + // caught up in the millisecond that it joined + late(5, start + 70_000, start + 70_000, 10, start + 80_000)); long[] epochs = {start + 1_000, start + 19_000, start + 39_000}; RunReport.Samples samples = new RunReport.Samples(epochs, new double[3], Map.of(), Map.of("app-1", new double[] {0, 570_000, 0}, "app-2", new double[] {0, 0, 1_170_000}, @@ -90,10 +92,11 @@ public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception // the threshold that the application wrote .contains("within 500 ms of its publishing") .contains("| `app-1` | 20.0 s | 570,000 | 10.0 s | 150,000 msg/s | 100.0 s |") - .contains("| `app-2` | 40.0 s | 1,170,000 | 90.0 s, after the gateways finished | 44,444 msg/s" + .contains("| `app-2` | 40.0 s | 1,170,000 | 80.4 s, after the gateways finished | 49,751 msg/s" + " | 100.0 s |") .contains("| `app-3` | 0.5 s | – | not caught up | 30,769 msg/s until its last message | 149.5 s |") .contains("| `app-4` | 60.0 s | – | not caught up | – | – |") + .contains("| `app-5` | 70.0 s | – | 0.0 s | – | 10.0 s |") .doesNotContain("`app-0`"); // without topic stats StringBuilder withoutSamples = new StringBuilder(); diff --git a/tests/performance/scenarios/docs/iot-telemetry.md b/tests/performance/scenarios/docs/iot-telemetry.md index 8b507044433cc..937501653a010 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -48,7 +48,8 @@ application-visible order across Key_Shared hash-range reassignment. the 5 applications join after the measurement starts, 3 of them 2 s apart, at 20, 22 and 24 s, and the last one at 60 s. Each late application has its subscription from the start, so its backlog builds up until it joins, and it then reads the backlog while the gateways keep publishing. The report's Catch-up section shows each application's - backlog when it joined and how long it took to catch up; the last one may only catch up after the gateways finish. + backlog when it joined and how long it took to catch up. An application that is still behind when the gateways + finish can't catch up, since no newer message comes; the report then shows when it received its last message. - **The broker's entry cache:** a subscription without consumers isn't an active cursor, so nothing keeps the entries for a late application before it joins. Once the 3 that join 2 s apart are reading, the entries that the first of them reads from storage are expected to be read by the others behind it, so they stay in the cache for diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java index 9e84d0e8a37f8..84a018c83c0a4 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java @@ -18,7 +18,8 @@ */ package org.apache.pulsar.tests.performance.tools; -import java.util.Set; +import java.util.Collections; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicLong; import java.util.function.LongSupplier; @@ -31,7 +32,8 @@ final class CatchUpTracker { private final int topicCount; private final long thresholdMillis; - private final Set caughtUpTopics = ConcurrentHashMap.newKeySet(); + // each topic's first receipt of a message within the threshold + private final Map caughtUpTopics = new ConcurrentHashMap<>(); private final AtomicLong joinEpochMs = new AtomicLong(); private final AtomicLong caughtUpEpochMs = new AtomicLong(); private final AtomicLong messagesWhenCaughtUp = new AtomicLong(); @@ -55,9 +57,12 @@ void received(String topic, long publishEpochMs, long receivedEpochMs, LongSuppl || receivedEpochMs - publishEpochMs > thresholdMillis) { return; } - if (caughtUpTopics.add(topic) && caughtUpTopics.size() >= topicCount - && caughtUpEpochMs.compareAndSet(0, receivedEpochMs)) { - messagesWhenCaughtUp.set(messages.getAsLong()); + if (caughtUpTopics.putIfAbsent(topic, receivedEpochMs) == null && caughtUpTopics.size() >= topicCount) { + // the last topic's receipt, which another listener may have recorded after this one + long caughtUp = Collections.max(caughtUpTopics.values()); + if (caughtUpEpochMs.compareAndSet(0, caughtUp)) { + messagesWhenCaughtUp.set(messages.getAsLong()); + } } } 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 2559b7c9460bf..7285e8e346486 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 @@ -204,8 +204,8 @@ boolean enabled() { long warmupTotalSeconds = minimumRuntimeSeconds - durationSeconds; require(applications.joinSeconds().stream().allMatch(seconds -> seconds != null && seconds >= 0 && warmupTotalSeconds + seconds < timeoutSeconds), - "applications.joinSeconds must be at least 0, and the warmup (" + warmupTotalSeconds - + " s) and the latest join must be within timeoutSeconds"); + "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"); if (timeoutSeconds < minimumRuntimeSeconds) { 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/PerformanceTool.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java index c8f26fe580084..0de34684f1574 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/PerformanceTool.java @@ -72,7 +72,8 @@ abstract static class ScenarioCommand implements Callable { Integer controlPort; @Option(names = "--run-id", - description = "Shared correlation ID, required when warmup is enabled; use a new ID per run") + description = "Shared correlation ID, required when warmup is enabled or applications join" + + " later; use a new ID per run") String runId; IotScenario scenario() throws Exception { diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java index 4062768834a6b..cac308db7f01e 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java @@ -354,7 +354,8 @@ void openPods(PulsarClientSharedResources sharedResources, AtomicLong opened) th synchronized (pods) { // a late application's joiner can open a pod after the application stopped if (stopping.get()) { - created.close(); + // uninterruptibly, so that the client is closed before the shared resources are + created.closeAsync().exceptionally(failure -> null).join(); return; } pods.add(created); @@ -423,30 +424,27 @@ boolean finish() throws Exception { /** Closes the application's pods, when it hasn't finished, such as after a failure. */ void close() throws Exception { stopping.set(true); - Thread joining = joiner; - if (joining != null) { - joining.interrupt(); - // so that it doesn't open pods after they're closed, or use the shared resources after they are - joining.join(TimeUnit.SECONDS.toMillis(10)); - } - Thread restarting = restarter; - if (restarting != null) { - restarting.interrupt(); - } + // so that they don't open pods after they're closed, or use the shared resources after they are + interruptAndAwait(joiner); + interruptAndAwait(restarter); closePods(); } private void stopRestarts() throws InterruptedException { stopping.set(true); - Thread joining = joiner; - if (joining != null) { - joining.interrupt(); - joining.join(TimeUnit.SECONDS.toMillis(10)); + interruptAndAwait(joiner); + interruptAndAwait(restarter); + } + + /** Interrupts the thread and waits for it, which takes up to a client's close timeout when it opens a pod. */ + private static void interruptAndAwait(Thread thread) throws InterruptedException { + if (thread == null) { + return; } - Thread restarting = restarter; - if (restarting != null) { - restarting.interrupt(); - restarting.join(TimeUnit.SECONDS.toMillis(10)); + thread.interrupt(); + thread.join(TimeUnit.SECONDS.toMillis(90)); + if (thread.isAlive()) { + System.err.println("WARN " + thread.getName() + " didn't stop within 90 s"); } } @@ -492,14 +490,17 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, (current, received) -> current == 0 ? received : Math.min(current, received)); lastMeasurementReceiptEpochMs.accumulateAndGet(receivedEpochMs, Math::max); - catchUp.received(message.getTopicName(), message.getPublishTime(), - receivedEpochMs, tracker::uniqueMessages); } receiveLatency.recordMillis(receivedEpochMs - message.getPublishTime(), decoded.measurement()); tracker.received(decoded.deviceId(), decoded.sequence(), message.getMessageId(), decoded.sentNanos(), message.getTopicName(), Thread.currentThread().getName()); + if (decoded.measurement()) { + // after counting the message, which then counts for the catch-up + catchUp.received(message.getTopicName(), message.getPublishTime(), + receivedEpochMs, tracker::uniqueMessages); + } currentConsumer.acknowledgeAsync(message); } catch (RuntimeException error) { tracker.invalidMessage(); @@ -509,7 +510,8 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, .subscribe(); return new ClientAndConsumer(client, consumer); } catch (Throwable error) { - client.close(); + // uninterruptibly: an interrupted subscription would otherwise leave the client closing + client.closeAsync().exceptionally(failure -> null).join(); throw error; } } @@ -541,8 +543,11 @@ private void restartPods(PulsarClientSharedResources sharedResources) { private record ClientAndConsumer(PulsarClient client, Consumer consumer) implements AutoCloseable { @Override public void close() throws Exception { - consumer.close(); - client.close(); + try { + consumer.close(); + } finally { + client.close(); + } } /** Closes the consumer, then the client, also when closing the consumer failed. */ diff --git a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java index c95c860454ea4..68dabacb601c8 100644 --- a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java @@ -45,6 +45,16 @@ public void catchesUpWhenEveryTopicDeliveredWithinTheThresholdAfterJoining() { assertThat(tracker.thresholdMillis()).isEqualTo(1_000); } + @Test + public void catchesUpAtTheLatestTopicsReceiptWhenListenersRecordOutOfOrder() { + CatchUpTracker tracker = new CatchUpTracker(2, 1_000); + tracker.joined(20_000); + // concurrent listeners: t1's receipt at 23,000 is recorded before t0's earlier one at 22,000 + tracker.received("t1", 22_900, 23_000, () -> 10); + tracker.received("t0", 21_900, 22_000, () -> 11); + assertThat(tracker.caughtUpEpochMs()).isEqualTo(23_000); + } + @Test public void anApplicationThatJoinedAtTheStartDoesntCatchUp() { CatchUpTracker tracker = new CatchUpTracker(1, 1_000); 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 ad1ce1dfedeec..a3c54820743c2 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 @@ -148,6 +148,13 @@ public void rejectsInvalidJoinSettings() { null, Arrays.asList(0, null), null))).hasMessageContaining("applications.joinSeconds must be"); assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), null, List.of(0, 20), 0))).hasMessageContaining("caughtUpLatencyMillis must be at least 1"); + // the joins count from the measurement's start, after 2 warmup rounds of 20 s with 2 s after each + IotScenario.Warmup warmup = new IotScenario.Warmup(20, 0, 2, 2); + assertThat(scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), null, + List.of(0, 255), null), warmup).joinSeconds(1)).isEqualTo(255); + assertThatThrownBy(() -> scenario(new IotScenario.Applications(2, 2, "app-", new IotScenario.Client(2, 2), + null, List.of(0, 256), null), warmup)) + .hasMessageContaining("the warmup (44 s) and the latest join must be within timeoutSeconds"); } @Test @@ -183,8 +190,13 @@ private static IotScenario scenario(int warmupSeconds, long warmupMessages, int } private static IotScenario scenario(IotScenario.Applications applications) { - return 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), + return scenario(applications, new IotScenario.Warmup(0, 0, 1, 0)); + } + + private static IotScenario scenario(IotScenario.Applications applications, IotScenario.Warmup warmup) { + // 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.Topics(2, "persistent://public/default/iot-"), applications, null, 300); } From 36e59e1f1ddf32565a64101338e9230f08efa510 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sun, 4 Oct 2026 03:21:47 +0300 Subject: [PATCH 5/5] [improve][test] Address the third catch-up review: quiet join timeout, catch-up result published together, docs Assisted-by: Claude Code (claude-opus-5-5) --- tests/performance/docs/run-reports.md | 6 ++++-- tests/performance/scenarios/README.md | 2 +- .../scenarios/docs/iot-telemetry.md | 4 ++-- .../performance/tools/CatchUpTracker.java | 21 +++++++++++-------- .../tools/MeasurementStartMarker.java | 9 ++++++-- .../performance/tools/TelemetryConsumer.java | 21 +++++++++++++++++++ .../tools/MeasurementStartMarkerTest.java | 2 ++ 7 files changed, 49 insertions(+), 16 deletions(-) diff --git a/tests/performance/docs/run-reports.md b/tests/performance/docs/run-reports.md index 9b7927cc87092..b4bd88e599e5b 100644 --- a/tests/performance/docs/run-reports.md +++ b/tests/performance/docs/run-reports.md @@ -147,7 +147,9 @@ report has these sections: application's join time, its subscription's backlog then, how long it took to catch up, its catch-up rate and when it received its last message. An application that didn't catch up shows its rate until its last message instead. When applications join later, the Throughput section's delivered throughput is mostly the last one's. The backlog chart - shows each subscription's backlog building up until its application joins and falling as it catches up. + shows each subscription's backlog building up until its application joins and falling as it catches up. A late + application's end-to-end latency includes the time its messages waited in the backlog, so its lines dominate the + Latency section's charts. - **Latency**: the publish latency (send to acknowledgment) and each application's end-to-end latency (publish to consume) at percentiles from p50 to the maximum, with charts by percentile and over time. - **Backlog and rates**: each subscription's backlog and the per-second rates, sampled from the broker's topic @@ -196,7 +198,7 @@ beside it, rendered with [commonmark-java](https://github.com/commonmark/commonm | `gateways/gateways-summary.json` | The gateways' counts and throughput, and the epoch-millisecond boundaries of the measurement | | `gateways/gateways-latency.hdr`, `.hgrm` | The publish latency log, and its percentile distribution in milliseconds, see [Latency logs](#latency-logs) | | `gateways/gateways-state.bin` | The gateways' next sequence number for each device. The launcher compares it with each application's `application-state.bin` and fails the run when they differ, which catches messages missing at the end, where no gap shows | -| `applications//application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt. Its join setting (`joinSeconds`), and for an application that joined after the measurement started, its join time (`joinEpochMs`), when it caught up (`caughtUpEpochMs`) within `caughtUpLatencyMillis`, and its received messages then (`messagesWhenCaughtUp`); the times are 0 for an application that joined at the start or didn't catch up | +| `applications//application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt. Its join setting (`joinSeconds`), and for an application that joined after the measurement started, its join time (`joinEpochMs`), when it caught up (`caughtUpEpochMs`) within `caughtUpLatencyMillis`, and its received messages then (`messagesWhenCaughtUp`); `joinEpochMs` is 0 for an application that joined at the start, and `caughtUpEpochMs` and `messagesWhenCaughtUp` are 0 for one that didn't catch up | | `applications//application-latency.hdr`, `.hgrm` | The application's end-to-end latency log, and its percentile distribution in milliseconds | | `applications//application-state.bin` | The application's next expected sequence number for each device | | `applications//ordering-violations.txt` | Samples of the ordering violations, with the message ID, topic and receiving thread; empty in a valid run | diff --git a/tests/performance/scenarios/README.md b/tests/performance/scenarios/README.md index e3fe215ed5c52..ab9eaf6e69c00 100644 --- a/tests/performance/scenarios/README.md +++ b/tests/performance/scenarios/README.md @@ -47,7 +47,7 @@ detail. | [`iot-telemetry-small-restarts.yaml`](iot-telemetry-small-restarts.yaml) | The smaller topology, restarting 10 % of each application's clients every 30 s | low, about 3 GB | | [`iot-telemetry-restarts.yaml`](iot-telemetry-restarts.yaml) | The full topology, restarting 10 % of each application's clients every 30 s | medium, about 11 GB | | [`iot-telemetry-high-rate.yaml`](iot-telemetry-high-rate.yaml) | A high rate: five million messages at 30,000 messages per second, from 500 preconnected producers to one topic, consumed by 5 applications with 10 pods each | high, about 14 GB | -| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 180 s, with 4 of its 5 applications joining 20, 22, 24 and 60 s after the measurement starts and catching up on their backlogs | high, about 14 GB | +| [`iot-telemetry-catch-up.yaml`](iot-telemetry-catch-up.yaml) | Consumers joining a live stream: the high-rate scenario's 30,000 messages per second for 180 s, with 4 of its 5 applications joining 20, 22, 24 and 60 s after the measurement starts and reading their backlogs | high, about 14 GB | | [`iot-telemetry-max-rate.yaml`](iot-telemetry-max-rate.yaml) | The maximum rate: the 500 gateways of the high-rate scenario publish four million unbatched messages to one topic without a rate limit, consumed by one application with 20 pods, on single-copy ledgers | high, about 14 GB | ## Configurations diff --git a/tests/performance/scenarios/docs/iot-telemetry.md b/tests/performance/scenarios/docs/iot-telemetry.md index 937501653a010..24a77bcdc50d7 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -178,8 +178,8 @@ workloads: listenerThreads: 16 env: # the applications' container, which runs every application; from the memory configuration PULSAR_MEM: -Xms1536m -Xmx1536m -XX:MaxDirectMemorySize=256m -XX:+UseTransparentHugePages -XX:+AlwaysPreTouch - joinSeconds: [] # each application's join, in seconds after the measurement starts; empty: all at 0 - caughtUpLatencyMillis: 1000 # a late application has caught up within this end-to-end latency + joinSeconds: [] # not set in the base file: each application's join after the measurement starts, in s + caughtUpLatencyMillis: 1000 # not set in the base file: a late application has caught up within this latency behaviors: podRestarts: # each application restarts this fraction of its pods every intervalSeconds intervalSeconds: 0 diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java index 84a018c83c0a4..acf7e86dc3fad 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java @@ -22,6 +22,7 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.LongSupplier; /** @@ -35,8 +36,11 @@ final class CatchUpTracker { // each topic's first receipt of a message within the threshold private final Map caughtUpTopics = new ConcurrentHashMap<>(); private final AtomicLong joinEpochMs = new AtomicLong(); - private final AtomicLong caughtUpEpochMs = new AtomicLong(); - private final AtomicLong messagesWhenCaughtUp = new AtomicLong(); + // when it caught up and its received messages then, published together + private final AtomicReference caughtUp = new AtomicReference<>(); + + private record Result(long epochMs, long messages) { + } CatchUpTracker(int topicCount, long thresholdMillis) { this.topicCount = topicCount; @@ -53,16 +57,13 @@ void joined(long epochMs) { * @param messages the application's received messages, read when it has caught up */ void received(String topic, long publishEpochMs, long receivedEpochMs, LongSupplier messages) { - if (joinEpochMs.get() == 0 || caughtUpEpochMs.get() != 0 + if (joinEpochMs.get() == 0 || caughtUp.get() != null || receivedEpochMs - publishEpochMs > thresholdMillis) { return; } if (caughtUpTopics.putIfAbsent(topic, receivedEpochMs) == null && caughtUpTopics.size() >= topicCount) { // the last topic's receipt, which another listener may have recorded after this one - long caughtUp = Collections.max(caughtUpTopics.values()); - if (caughtUpEpochMs.compareAndSet(0, caughtUp)) { - messagesWhenCaughtUp.set(messages.getAsLong()); - } + caughtUp.compareAndSet(null, new Result(Collections.max(caughtUpTopics.values()), messages.getAsLong())); } } @@ -71,11 +72,13 @@ long joinEpochMs() { } long caughtUpEpochMs() { - return caughtUpEpochMs.get(); + Result result = caughtUp.get(); + return result != null ? result.epochMs() : 0; } long messagesWhenCaughtUp() { - return messagesWhenCaughtUp.get(); + Result result = caughtUp.get(); + return result != null ? result.messages() : 0; } long thresholdMillis() { diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java index f66a43cc71a4e..e98139c50c1e6 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java @@ -24,6 +24,7 @@ import java.nio.file.NoSuchFileException; import java.nio.file.Path; import java.nio.file.StandardCopyOption; +import java.util.concurrent.TimeoutException; /** * The gateways' measurement start, for the applications that join after it: a file in the coordination directory @@ -43,7 +44,11 @@ static void mark(Path directory, String runId, long epochMs) throws IOException Files.move(temporary, marker, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); } - /** Waits for the measurement to start, and returns its epoch milliseconds. */ + /** + * Waits for the measurement to start, and returns its epoch milliseconds. + * + * @throws TimeoutException when the measurement doesn't start before the deadline + */ static long await(Path directory, String runId, long deadlineNanos) throws Exception { Path marker = marker(directory, runId); while (deadlineNanos - System.nanoTime() > 0) { @@ -53,7 +58,7 @@ static long await(Path directory, String runId, long deadlineNanos) throws Excep Thread.sleep(POLL_INTERVAL_MILLIS); } } - throw new IllegalStateException("Timed out waiting for the measurement to start"); + throw new TimeoutException("Timed out waiting for the measurement to start"); } private static Path marker(Path directory, String runId) { diff --git a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java index cac308db7f01e..c6e0d613a85ca 100644 --- a/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/TelemetryConsumer.java @@ -37,7 +37,9 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import org.apache.pulsar.client.api.Consumer; @@ -260,6 +262,7 @@ private static final class Application { // For an application that joins later: when it joined and when it caught up private final CatchUpTracker catchUp; private final AtomicReference joinFailure = new AtomicReference<>(); + private final AtomicInteger listenersInFlight = new AtomicInteger(); private volatile Thread joiner; // Guarded by itself, as the restarts replace pods private final List pods; @@ -334,6 +337,9 @@ void scheduleJoin(PulsarClientSharedResources sharedResources, Path coordination startRestarts(sharedResources); } catch (InterruptedException interrupted) { Thread.currentThread().interrupt(); + } catch (TimeoutException measurementDidNotStart) { + // the receive loop's timeout reports the run as failed and writes the applications' summaries + System.out.println("JOIN application=" + index + " didn't join: the measurement didn't start"); } catch (Throwable error) { joinFailure.compareAndSet(null, error); } @@ -397,6 +403,7 @@ boolean receivedEveryMessage() { */ boolean finish() throws Exception { stopRestarts(); + awaitListeners(); DeviceSequenceTracker.Summary summary = tracker.summary(); receiveLatency.close(); tracker.writeState(output.resolve("application-state.bin")); @@ -436,6 +443,17 @@ private void stopRestarts() throws InterruptedException { interruptAndAwait(restarter); } + /** + * Waits briefly for the listeners that are handling a message, such as the one that received the last message + * and may record the catch-up with it, so that the summary includes what they record. + */ + private void awaitListeners() throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(1); + while (listenersInFlight.get() > 0 && deadline - System.nanoTime() > 0) { + Thread.sleep(1); + } + } + /** Interrupts the thread and waits for it, which takes up to a client's close timeout when it opens a pod. */ private static void interruptAndAwait(Thread thread) throws InterruptedException { if (thread == null) { @@ -477,6 +495,7 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) .messageListener((currentConsumer, message) -> { long receivedEpochMs = System.currentTimeMillis(); + listenersInFlight.incrementAndGet(); try { TelemetryMessage.Decoded decoded = TelemetryMessage.decode(message.getData()); byte[] key = message.getKeyBytes(); @@ -505,6 +524,8 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, } catch (RuntimeException error) { tracker.invalidMessage(); currentConsumer.negativeAcknowledge(message); + } finally { + listenersInFlight.decrementAndGet(); } }) .subscribe(); diff --git a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java index 23e575091139e..55a62fe8e2ea3 100644 --- a/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java @@ -24,6 +24,7 @@ import java.nio.file.Path; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import org.testng.annotations.Test; public class MeasurementStartMarkerTest { @@ -45,6 +46,7 @@ public void applicationsAwaitTheGatewaysMeasurementStart() throws Exception { // another run's marker isn't this run's assertThatThrownBy(() -> MeasurementStartMarker.await(directory, "run-2", System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(100))) + .isInstanceOf(TimeoutException.class) .hasMessageContaining("Timed out waiting for the measurement to start"); } }