Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion tests/performance/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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>/application-summary.json` | Each application's unique messages, duplicates, ordering violations and invalid messages; `lastMeasurementMessageReceivedEpochMs` records the last measured receipt | `jq` |
| `applications/<application>/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/<application>/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>/application-latency.hgrm` | The publish and end-to-end latency percentile distributions, in milliseconds, as text | Read |
Expand Down
13 changes: 10 additions & 3 deletions tests/performance/docs/run-reports.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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>/application-summary.json` | The application's unique messages, duplicates, ordering violations and invalid messages, and its first and last measured-message receipt |
| `applications/<application>/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>/application-latency.hdr`, `.hgrm` | The application's end-to-end latency log, and its percentile distribution in milliseconds |
| `applications/<application>/application-state.bin` | The application's next expected sequence number for each device |
| `applications/<application>/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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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<JsonNode> consumers, Samples samples,
long measurementStart, long measurementEnd) {
List<JsonNode> 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)
: "–"));
}
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<JsonNode> 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"));
Expand Down
1 change: 1 addition & 0 deletions tests/performance/scenarios/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading