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
30 changes: 22 additions & 8 deletions benchmarks/src/main/java/io/hstore/bench/Comparison.java
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ record Result(String workload, String description, long operations, double milli
this(workload, description, operations, millis, checksum, Map.of());
}

Result using(JvmUsage usage) {
return with("gcCount", number(usage.collections())).with("gcMillis", number(usage.pauseMillis()))
.with("allocated", number(usage.allocated()));
}

Result with(String key, Json value) {
Map<String, Json> fields = new LinkedHashMap<>(extra);
fields.put(key, value);
Expand Down Expand Up @@ -94,7 +99,8 @@ private static void run(Map<String, String> options) throws Exception {
Path directory = Files.createTempDirectory("hstore-bench-" + storeName);
Dataset dataset = Dataset.generate(scale, SEED);
List<Result> results = new ArrayList<>();
long written = 0;
Map<String, Json> totals = new LinkedHashMap<>();
long baselineHeap = JvmUsage.retainedHeap();
Store store = switch (storeName) {
case "hstore" -> new HStoreStore(directory, sync, Integer.parseInt(options.getOrDefault("history", "64")));
case "hypergraphdb" -> new HyperGraphDbStore(directory, sync);
Expand All @@ -112,6 +118,7 @@ private static void run(Map<String, String> options) throws Exception {
store.ingestEdges(dataset, BATCH);
return dataset.incidences();
})));
totals.put("retainedHeap", number(JvmUsage.retainedHeap() - baselineHeap));
int[] nodes = dataset.probeNodes();
results.add(read("read.incidence", "enumerate the incidence set of a node", nodes.length,
(from, to) -> store.incidence(nodes, from, to), true));
Expand Down Expand Up @@ -154,21 +161,26 @@ private static void run(Map<String, String> options) throws Exception {
results.add(read("read.incidence.cold", "incidence sets immediately after reopening", nodes.length,
(from, to) -> store.incidence(nodes, from, to), false));
store.flush();
written = store.bytesWritten();
totals.put("bytesWritten", number(store.bytesWritten()));
} finally {
store.close();
}
long disk = store.diskBytes();
results.add(new Result("disk", "bytes on disk after a clean shutdown", disk, 1000.0, disk));
write(out, store, scale, sync, threads, dataset, written, results);
write(out, store, scale, sync, threads, dataset, totals, results);
results.forEach(result -> IO.println(" %-22s %,14.0f ops/s %,10.1f ms checksum %d"
.formatted(result.workload(), result.rate(), result.millis(), result.checksum())));
}

private static Json number(long value) {
return new Json.Number(BigDecimal.valueOf(value));
}

private static Result measure(String workload, String description, long operations, LongSupplier work) {
JvmUsage usage = JvmUsage.now();
long started = System.nanoTime();
long checksum = work.getAsLong();
return new Result(workload, description, operations, (System.nanoTime() - started) / 1e6, checksum);
return new Result(workload, description, operations, (System.nanoTime() - started) / 1e6, checksum).using(JvmUsage.now().since(usage));
}

private static long batched(int length, Slice slice) {
Expand Down Expand Up @@ -201,19 +213,21 @@ private static Result latency(String workload, String description, int operation
}
long[] nanos = new long[operations];
long checksum = 0;
JvmUsage usage = JvmUsage.now();
long started = System.nanoTime();
for (int i = 0; i < operations; i++) {
long begin = System.nanoTime();
checksum += operation.applyAsLong(i);
nanos[i] = System.nanoTime() - begin;
}
return new Result(workload, description, operations, (System.nanoTime() - started) / 1e6, checksum)
.with("latency", Latency.of(nanos).json());
.using(JvmUsage.now().since(usage)).with("latency", Latency.of(nanos).json());
}

private static Result concurrent(Store store, Dataset dataset, int threads) throws Exception {
int[] probes = dataset.probeNodes();
try (ExecutorService pool = Executors.newFixedThreadPool(threads)) {
JvmUsage usage = JvmUsage.now();
long started = System.nanoTime();
List<Future<Long>> futures = new ArrayList<>();
for (int t = 0; t < threads; t++) {
Expand All @@ -230,11 +244,11 @@ private static Result concurrent(Store store, Dataset dataset, int threads) thro
}
double millis = (System.nanoTime() - started) / 1e6;
return new Result("read.incidence.parallel", "incidence sets from %d threads, %,d probes each".formatted(threads, probes.length),
(long) threads * probes.length, millis, checksum);
(long) threads * probes.length, millis, checksum).using(JvmUsage.now().since(usage));
}
}

private static void write(Path out, Store store, int scale, boolean sync, int threads, Dataset dataset, long written, List<Result> results) {
private static void write(Path out, Store store, int scale, boolean sync, int threads, Dataset dataset, Map<String, Json> totals, List<Result> results) {
Map<String, Json> fields = new LinkedHashMap<>();
fields.put("store", new Json.Str(store.name()));
fields.put("version", new Json.Str(store.version()));
Expand All @@ -249,7 +263,7 @@ private static void write(Path out, Store store, int scale, boolean sync, int th
fields.put("os", new Json.Str(System.getProperty("os.name") + " " + System.getProperty("os.version") + " " + System.getProperty("os.arch")));
fields.put("processors", new Json.Number(BigDecimal.valueOf(Runtime.getRuntime().availableProcessors())));
fields.put("maxHeap", new Json.Number(BigDecimal.valueOf(Runtime.getRuntime().maxMemory())));
fields.put("bytesWritten", new Json.Number(BigDecimal.valueOf(written)));
fields.putAll(totals);
fields.put("results", new Json.Array(results.stream().map(Result::json).toList()));
try {
Files.createDirectories(out.toAbsolutePath().getParent());
Expand Down
28 changes: 28 additions & 0 deletions benchmarks/src/main/java/io/hstore/bench/JvmUsage.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
package io.hstore.bench;

import java.lang.management.GarbageCollectorMXBean;
import java.lang.management.ManagementFactory;

record JvmUsage(long collections, long pauseMillis, long allocated) {

static JvmUsage now() {
long collections = 0;
long pauseMillis = 0;
for (GarbageCollectorMXBean collector : ManagementFactory.getGarbageCollectorMXBeans()) {
collections += Math.max(0, collector.getCollectionCount());
pauseMillis += Math.max(0, collector.getCollectionTime());
}
long allocated = ((com.sun.management.ThreadMXBean) ManagementFactory.getThreadMXBean()).getTotalThreadAllocatedBytes();
return new JvmUsage(collections, pauseMillis, allocated);
}

JvmUsage since(JvmUsage start) {
return new JvmUsage(collections - start.collections, pauseMillis - start.pauseMillis, allocated - start.allocated);
}

static long retainedHeap() {
System.gc();
System.gc();
return ManagementFactory.getMemoryMXBean().getHeapMemoryUsage().getUsed();
}
}
21 changes: 16 additions & 5 deletions benchmarks/src/main/java/io/hstore/bench/Report.java
Original file line number Diff line number Diff line change
Expand Up @@ -73,18 +73,24 @@ static String markdown(String[] args) {
private record Metric(String key, String name, boolean perOperation, DoubleFunction<String> format) {
}

private static final List<Metric> METRICS = List.of(new Metric("bytesWritten", "bytes written per operation", true, Report::bytes));
private static final List<Metric> METRICS = List.of(new Metric("bytesWritten", "bytes written per operation", true, Report::bytes),
new Metric("allocated", "heap allocated per operation", true, Report::bytes),
new Metric("gcMillis", "GC pause time", false, value -> "%,.0f ms".formatted(value)));

private static void resources(StringBuilder out, List<Run> subjects, List<Run> baselines) {
List<Run> all = Stream.concat(subjects.stream(), baselines.stream()).toList();
StringBuilder rows = new StringBuilder();
for (String workload : subjects.getFirst().results().keySet()) {
for (Metric metric : METRICS) {
for (Metric metric : METRICS) {
for (String workload : subjects.getFirst().results().keySet()) {
if (all.stream().allMatch(run -> run.results().containsKey(workload) && run.results().get(workload).containsKey(metric.key()))) {
double mine = median(metric(subjects, workload, metric));
double theirs = median(metric(baselines, workload, metric));
rows.append("| `%s` | %s | %s | %s | **%.2f×** |%n".formatted(workload, metric.name(), metric.format().apply(mine),
metric.format().apply(theirs), theirs / mine));
if (mine == 0 && theirs == 0) {
continue;
}
double ratio = theirs / mine;
rows.append("| `%s` | %s | %s | %s | %s |%n".formatted(workload, metric.name(), metric.format().apply(mine),
metric.format().apply(theirs), Double.isFinite(ratio) ? "**%.2f×**".formatted(ratio) : "n/a"));
}
}
}
Expand All @@ -94,6 +100,11 @@ private static void resources(StringBuilder out, List<Run> subjects, List<Run> b
out.append("%nResources used by each workload, median of the runs:%n%n".formatted());
out.append("| Workload | Measure | %s | %s | Ratio |%n".formatted(subjects.getFirst().store(), baselines.getFirst().store()));
out.append("|---|---|---:|---:|---:|%n".formatted()).append(rows);
if (all.stream().allMatch(run -> run.header().containsKey("retainedHeap"))) {
out.append("%nHeap retained by the loaded store after a full GC: %s %s, %s %s.%n".formatted(
subjects.getFirst().store(), bytes(median(header(subjects, "retainedHeap"))),
baselines.getFirst().store(), bytes(median(header(baselines, "retainedHeap")))));
}
if (all.stream().allMatch(run -> run.header().containsKey("bytesWritten"))) {
out.append("%nBytes written in total, after a final flush: %s %s, %s %s.%n".formatted(
subjects.getFirst().store(), bytes(median(header(subjects, "bytesWritten"))),
Expand Down
7 changes: 7 additions & 0 deletions docs/benchmarks.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,13 @@ so bytes still sitting in buffers are counted against the workload that produced
when a store reopens, so the adapters carry the total across `reopen`. Bytes written during the final close
aren't counted.

**Memory.** Every workload also records, from the JVM's own MXBeans, the bytes the whole process allocated
(`getTotalThreadAllocatedBytes`) and the GC pause time (`GarbageCollectorMXBean.getCollectionTime`, which is
stop-the-world time with ParallelGC). Rows where both stores paused for 0 ms are left out. After ingest, the run
forces a full GC and records how much heap the loaded store keeps. A baseline taken before the store opens is
subtracted, so the shared dataset arrays don't count. This number mostly reflects how each engine's cache is
configured (see *Setup*).

## Results

Ratios above 1 favour HStore. Throughput ratios divide HStore by HyperGraphDB; latency and size ratios divide
Expand Down
Loading