diff --git a/tests/performance/AGENTS.md b/tests/performance/AGENTS.md index 557a4f41ac7d9..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 | `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 59a5a090d093a..b4bd88e599e5b 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, @@ -143,6 +143,13 @@ 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, 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. 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 @@ -191,12 +198,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. 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 | | `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 7d611339e98c1..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 @@ -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, measurementEnd); + 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) { @@ -640,6 +642,69 @@ 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 finish;" + + " 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, 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 measurementEnd) { + List late = consumers.stream().filter(consumer -> consumer.path("joinEpochMs").asLong() > 0) + .toList(); + if (late.isEmpty()) { + return; + } + 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 each of its topics delivered a measured message within ") + .append(String.format(Locale.ROOT, "%,d", caughtUpLatency)) + .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(); + 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; + } + } + } + String caughtUpAfter; + String rate; + if (caughtUp > 0 && caughtUp >= joined) { + caughtUpAfter = String.format(Locale.ROOT, "%.1f s", (caughtUp - joined) / 1000.0) + + (caughtUp > measurementEnd ? ", after the gateways finished" : ""); + 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", + 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, caughtUpAfter, rate, + lastReceived > joined ? String.format(Locale.ROOT, "%.1f s", (lastReceived - joined) / 1000.0) + : "–")); + } } /** 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..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 @@ -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,58 @@ public void deleteRun() throws IOException { } } + @Test + public void reportsTheCatchUpOfTheApplicationsThatJoinedLater() throws Exception { + 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}"), + // caught up while the gateways published + late(1, start + 20_000, start + 30_000, 1_500_000, start + 120_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), + // 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}, + "app-3", new double[] {0, 0, 0})); + StringBuilder report = new StringBuilder(); + 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") + .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 | 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(); + 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, 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 35d0fbc32a355..ab9eaf6e69c00 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 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 376225662a7de..24a77bcdc50d7 100644 --- a/tests/performance/scenarios/docs/iot-telemetry.md +++ b/tests/performance/scenarios/docs/iot-telemetry.md @@ -43,6 +43,22 @@ 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 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. 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 + 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 @@ -162,12 +178,24 @@ 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: [] # 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 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 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 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..d68eeabed9d23 --- /dev/null +++ b/tests/performance/scenarios/iot-telemetry-catch-up.yaml @@ -0,0 +1,33 @@ +# 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. 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 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: 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..acf7e86dc3fad --- /dev/null +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/CatchUpTracker.java @@ -0,0 +1,87 @@ +/* + * 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.Collections; +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; + +/** + * 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; + // each topic's first receipt of a message within the threshold + private final Map caughtUpTopics = new ConcurrentHashMap<>(); + private final AtomicLong joinEpochMs = 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; + 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 || 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 + caughtUp.compareAndSet(null, new Result(Collections.max(caughtUpTopics.values()), messages.getAsLong())); + } + } + + long joinEpochMs() { + return joinEpochMs.get(); + } + + long caughtUpEpochMs() { + Result result = caughtUp.get(); + return result != null ? result.epochMs() : 0; + } + + long messagesWhenCaughtUp() { + Result result = caughtUp.get(); + return result != null ? result.messages() : 0; + } + + 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 b3e263f9daefd..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 @@ -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; @@ -94,9 +95,29 @@ 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 { + // 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; + } + } + + 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,8 +197,17 @@ 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"); + // 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 + && warmupTotalSeconds + seconds < 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) { - 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" @@ -239,6 +269,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..e98139c50c1e6 --- /dev/null +++ b/tests/performance/tools/src/main/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarker.java @@ -0,0 +1,67 @@ +/* + * 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; +import java.util.concurrent.TimeoutException; + +/** + * 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. + * + * @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) { + try { + return Long.parseLong(Files.readString(marker, StandardCharsets.UTF_8).trim()); + } catch (NoSuchFileException e) { + Thread.sleep(POLL_INTERVAL_MILLIS); + } + } + throw new TimeoutException("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..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 { @@ -81,6 +82,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..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; @@ -56,6 +58,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 +96,24 @@ 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"); + // 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) { - 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 +149,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(); @@ -142,6 +163,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; @@ -154,7 +176,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 +197,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,11 +259,16 @@ 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 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; 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; @@ -245,6 +277,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); @@ -254,10 +287,83 @@ 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); + } + // 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) { + 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); + } + }, "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); synchronized (pods) { + // a late application's joiner can open a pod after the application stopped + if (stopping.get()) { + // uninterruptibly, so that the client is closed before the shared resources are + created.closeAsync().exceptionally(failure -> null).join(); + return; + } pods.add(created); } opened.incrementAndGet(); @@ -265,9 +371,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(); } } @@ -296,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")); @@ -309,7 +417,12 @@ boolean finish() throws Exception { + " \"firstMeasurementMessageReceivedEpochMs\": " + firstMeasurementReceiptEpochMs.get() + ",\n" + " \"lastMeasurementMessageReceivedEpochMs\": " - + lastMeasurementReceiptEpochMs.get() + "\n}\n"); + + lastMeasurementReceiptEpochMs.get() + ",\n" + + " \"joinSeconds\": " + scenario.joinSeconds(index) + ",\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(); @@ -318,17 +431,38 @@ 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 (restarter != null) { - restarter.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); - if (restarter != null) { - restarter.interrupt(); - restarter.join(TimeUnit.SECONDS.toMillis(10)); + interruptAndAwait(joiner); + 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) { + return; + } + thread.interrupt(); + thread.join(TimeUnit.SECONDS.toMillis(90)); + if (thread.isAlive()) { + System.err.println("WARN " + thread.getName() + " didn't stop within 90 s"); } } @@ -361,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(); @@ -380,16 +515,24 @@ private ClientAndConsumer createPod(PulsarClientSharedResources sharedResources, 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(); currentConsumer.negativeAcknowledge(message); + } finally { + listenersInFlight.decrementAndGet(); } }) .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; } } @@ -421,8 +564,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/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/CatchUpTrackerTest.java b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java new file mode 100644 index 0000000000000..68dabacb601c8 --- /dev/null +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/CatchUpTrackerTest.java @@ -0,0 +1,64 @@ +/* + * 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 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); + 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 4f2ae2955cd07..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 @@ -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,65 @@ 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, + 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("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, 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 public void rejectsTimeBasedWarmupWithoutRateLimit() { assertThatThrownBy(() -> scenario(20, 0, 0, 5_000_000)) @@ -129,6 +189,18 @@ private static IotScenario scenario(int warmupSeconds, long warmupMessages, int 300); } + private static IotScenario scenario(IotScenario.Applications applications) { + 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); + } + 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..55a62fe8e2ea3 --- /dev/null +++ b/tests/performance/tools/src/test/java/org/apache/pulsar/tests/performance/tools/MeasurementStartMarkerTest.java @@ -0,0 +1,52 @@ +/* + * 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 java.util.concurrent.TimeoutException; +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))) + .isInstanceOf(TimeoutException.class) + .hasMessageContaining("Timed out waiting for the measurement to start"); + } +}