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
75 changes: 69 additions & 6 deletions benchmarks/src/main/java/io/hstore/bench/Comparison.java
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,10 @@

import io.hstore.db.value.Json;

import java.io.BufferedReader;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.lang.management.ManagementFactory;
import java.math.BigDecimal;
import java.nio.file.Files;
import java.nio.file.Path;
Expand Down Expand Up @@ -63,6 +65,7 @@ private interface Slice {

private static final int BATCH = 1_000;
private static final int COMMITS = 5_000;
private static final String COMMITTED = "committed ";
private static final long MIXED_NANOS = 10_000_000_000L;
private static final int MIXED_BATCH = 10;
private static final int MIXED_UPDATES_PER_SECOND = 1_000;
Expand All @@ -76,10 +79,12 @@ public static void main(String[] args) throws Exception {
Map<String, String> options = options(args);
switch (args.length == 0 ? "" : args[0]) {
case "run" -> run(options);
case "crash" -> crash(options);
case "report" -> IO.println(Report.markdown(args));
default -> {
IO.println("""
usage: Comparison run --store hstore|hypergraphdb [--scale N] [--durability async|sync] [--threads N] [--history N] [--cache-mb N] --out FILE
Comparison crash --store hstore|hypergraphdb --dir DIRECTORY (used by run; ingests until killed)
Comparison report FILE... (runs of two stores; cells show the median and range)""");
System.exit(2);
}
Expand All @@ -103,16 +108,11 @@ private static void run(Map<String, String> options) throws Exception {
int threads = Integer.parseInt(options.getOrDefault("threads", String.valueOf(Runtime.getRuntime().availableProcessors())));
Path out = Path.of(options.getOrDefault("out", "results/" + storeName + ".json"));
Path directory = Files.createTempDirectory("hstore-bench-" + storeName);
long cacheBytes = Long.parseLong(options.getOrDefault("cache-mb", "0")) << 20;
Dataset dataset = Dataset.generate(scale, SEED);
List<Result> results = new ArrayList<>();
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")), cacheBytes);
case "hypergraphdb" -> new HyperGraphDbStore(directory, sync, cacheBytes);
default -> throw new IllegalArgumentException("unknown store " + storeName);
};
Store store = open(options, directory, true);
IO.println("%s: %,d nodes, %,d hyperedges, %,d incidences, durability %s, %d threads"
.formatted(store.name(), dataset.nodes(), dataset.edges().length, dataset.incidences(), sync ? "sync" : "async", threads));
try {
Expand Down Expand Up @@ -196,11 +196,74 @@ private static void run(Map<String, String> options) throws Exception {
}
long churned = store.diskBytes();
results.add(new Result("disk.churned", "bytes on disk after the deletes and a clean shutdown", churned, 1000.0, churned));
results.add(recover(options, dataset));
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 Store open(Map<String, String> options, Path directory, boolean create) {
boolean sync = options.getOrDefault("durability", "async").equals("sync");
long cacheBytes = Long.parseLong(options.getOrDefault("cache-mb", "0")) << 20;
return switch (options.getOrDefault("store", "hstore")) {
case "hstore" -> new HStoreStore(directory, sync, Integer.parseInt(options.getOrDefault("history", "64")), cacheBytes, create);
case "hypergraphdb" -> new HyperGraphDbStore(directory, sync, cacheBytes);
case String other -> throw new IllegalArgumentException("unknown store " + other);
};
}

private static void crash(Map<String, String> options) throws InterruptedException {
Dataset dataset = Dataset.generate(Integer.parseInt(options.getOrDefault("scale", "1")), SEED);
Store store = open(options, Path.of(options.get("dir")), true);
store.ingestNodes(dataset, BATCH, committed -> {
System.out.println(COMMITTED + committed);
System.out.flush();
});
Thread.sleep(Long.MAX_VALUE);
}

private static Result recover(Map<String, String> options, Dataset dataset) throws Exception {
Path directory = Files.createTempDirectory("hstore-bench-crash");
List<String> command = new ArrayList<>();
command.add(ProcessHandle.current().info().command().orElse("java"));
command.addAll(ManagementFactory.getRuntimeMXBean().getInputArguments());
command.addAll(List.of("-cp", System.getProperty("java.class.path"), Comparison.class.getName(), "crash", "--dir", directory.toString()));
for (String option : List.of("store", "scale", "durability", "history", "cache-mb")) {
if (options.containsKey(option)) {
command.addAll(List.of("--" + option, options.get(option)));
}
}
Process child = new ProcessBuilder(command).redirectErrorStream(true).start();
int target = dataset.nodes() / 2;
int acknowledged = 0;
try (BufferedReader lines = child.inputReader()) {
for (String line = lines.readLine(); line != null && acknowledged < target; line = lines.readLine()) {
if (line.startsWith(COMMITTED)) {
acknowledged = Integer.parseInt(line.substring(COMMITTED.length()));
}
}
} finally {
child.destroyForcibly();
child.waitFor();
}
if (acknowledged < target) {
throw new IllegalStateException("the crash child exited after %,d acknowledged nodes".formatted(acknowledged));
}
Store[] recovered = new Store[1];
Result result = measure("recover", "open the store after killing the ingest process with SIGKILL", 1, () -> {
recovered[0] = open(options, directory, false);
return 1;
});
long present;
try {
present = recovered[0].countNodes();
} finally {
recovered[0].close();
}
return result.with("checked", new Json.Bool(false)).with("acknowledged", number(acknowledged))
.with("lost", number(Math.max(0, acknowledged - present)));
}

private static Json number(long value) {
return new Json.Number(BigDecimal.valueOf(value));
}
Expand Down
23 changes: 16 additions & 7 deletions benchmarks/src/main/java/io/hstore/bench/HStoreStore.java
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.function.IntConsumer;

final class HStoreStore implements Store {

Expand All @@ -37,7 +38,7 @@ final class HStoreStore implements Store {
private final int cachedNodes;
private final long cacheBytes;

HStoreStore(Path directory, boolean sync, int history, long cacheBytes) {
HStoreStore(Path directory, boolean sync, int history, long cacheBytes, boolean create) {
this.directory = directory;
this.history = history;
this.cacheBytes = cacheBytes;
Expand All @@ -46,11 +47,13 @@ final class HStoreStore implements Store {
this.options = DatabaseOptions.defaults().withEngine(engine
.withDurability(sync ? Durability.SYNC : Durability.ASYNC).withCachedNodes(cachedNodes).withHistoryLimit(history));
this.database = HypergraphDatabase.open(directory, options);
database.write(writer -> {
writer.defineNode(NODE, List.of(new PropertyDef("name", TypeTag.STRING, false, false)));
writer.defineEdge(EDGE, AtomKind.SET_EDGE, List.of(), List.of());
return null;
});
if (create) {
database.write(writer -> {
writer.defineNode(NODE, List.of(new PropertyDef("name", TypeTag.STRING, false, false)));
writer.defineEdge(EDGE, AtomKind.SET_EDGE, List.of(), List.of());
return null;
});
}
}

@Override
Expand All @@ -75,7 +78,7 @@ public Path directory() {
}

@Override
public void ingestNodes(Dataset dataset, int batch) {
public void ingestNodes(Dataset dataset, int batch, IntConsumer committed) {
nodes = new long[dataset.nodes()];
for (int start = 0; start < nodes.length; start += batch) {
int from = start;
Expand All @@ -86,9 +89,15 @@ public void ingestNodes(Dataset dataset, int batch) {
}
return null;
});
committed.accept(to);
}
}

@Override
public long countNodes() {
return database.read(reader -> reader.atoms(NODE).count());
}

@Override
public void ingestEdges(Dataset dataset, int batch) {
int[][] source = dataset.edges();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import java.util.HashSet;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.function.IntConsumer;

final class HyperGraphDbStore implements Store {

Expand Down Expand Up @@ -78,7 +79,7 @@ private <T> T read(Callable<T> work) {
}

@Override
public void ingestNodes(Dataset dataset, int batch) {
public void ingestNodes(Dataset dataset, int batch, IntConsumer committed) {
nodes = new HGPersistentHandle[dataset.nodes()];
for (int start = 0; start < nodes.length; start += batch) {
int from = start;
Expand All @@ -89,9 +90,15 @@ public void ingestNodes(Dataset dataset, int batch) {
}
return null;
});
committed.accept(to);
}
}

@Override
public long countNodes() {
return read(() -> hg.count(graph, hg.type(String.class)));
}

@Override
public void ingestEdges(Dataset dataset, int batch) {
int[][] source = dataset.edges();
Expand Down
11 changes: 6 additions & 5 deletions benchmarks/src/main/java/io/hstore/bench/Report.java
Original file line number Diff line number Diff line change
Expand Up @@ -71,12 +71,13 @@ static String markdown(String[] args) {
return out.toString();
}

private record Metric(String key, String name, boolean perOperation, DoubleFunction<String> format) {
private record Metric(String key, String name, boolean perOperation, boolean showZero, DoubleFunction<String> format) {
}

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 final List<Metric> METRICS = List.of(new Metric("bytesWritten", "bytes written per operation", true, false, Report::bytes),
new Metric("allocated", "heap allocated per operation", true, false, Report::bytes),
new Metric("gcMillis", "GC pause time", false, false, value -> "%,.0f ms".formatted(value)),
new Metric("lost", "acknowledged nodes lost", false, true, value -> "%,.0f".formatted(value)));

private static void resources(StringBuilder out, List<Run> subjects, List<Run> baselines) {
List<Run> all = Stream.concat(subjects.stream(), baselines.stream()).toList();
Expand All @@ -86,7 +87,7 @@ private static void resources(StringBuilder out, List<Run> subjects, List<Run> b
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));
if (mine == 0 && theirs == 0) {
if (mine == 0 && theirs == 0 && !metric.showZero()) {
continue;
}
double ratio = theirs / mine;
Expand Down
10 changes: 9 additions & 1 deletion benchmarks/src/main/java/io/hstore/bench/Store.java
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.function.IntConsumer;
import java.util.stream.Stream;

interface Store extends AutoCloseable {
Expand All @@ -14,7 +15,14 @@ interface Store extends AutoCloseable {

String cache();

void ingestNodes(Dataset dataset, int batch);
default void ingestNodes(Dataset dataset, int batch) {
ingestNodes(dataset, batch, _ -> {
});
}

void ingestNodes(Dataset dataset, int batch, IntConsumer committed);

long countNodes();

void ingestEdges(Dataset dataset, int batch);

Expand Down
13 changes: 13 additions & 0 deletions docs/benchmarks.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ taken from a common hyperedge (so co-membership counts are non-zero), and 20,000
| `churn.remove` | remove one member from each of 20,000 other hyperedges, 1,000 per transaction | `Writer.remove` | `graph.replace` with a new `HGPlainLink` (links are immutable) |
| `read.incidence.churned` | `read.incidence` after the deletes | | |
| `disk.churned` | bytes on disk after the deletes and a clean shutdown | | |
| `recover` | open the store after the process ingesting into it was killed with SIGKILL | `HypergraphDatabase.open` | `HGEnvironment.get` |

Read workloads run a warm-up pass over 10% of the probes before timing, except `read.incidence.cold`.

Expand Down Expand Up @@ -92,6 +93,18 @@ the loaded store after a clean shutdown. The checksum of `read.incidence.churned
both engines applied the same deletes. `disk.churned` shows whether the space came back: through compaction for
HStore, and the log cleaner for JE.

**Crash recovery.** `reopen` closes the store cleanly first. `recover` doesn't:
1. The run starts a second JVM (`Comparison crash`) with the same flags and classpath. It ingests nodes into a
fresh directory and prints the count after every acknowledged commit.
2. Once half the nodes are acknowledged, the parent kills it with SIGKILL.
3. The parent times opening the directory, then counts the nodes that survived.

The *Resources* table shows how many acknowledged nodes were lost, including when it's zero. SIGKILL is a process
crash, not a power failure: bytes the engine handed to the operating system survive it even without an fsync. So
this measures recovery time and the engine's own buffering, not whether `sync` really reaches the disk. Torn
writes and lost fsyncs are covered by HStore's crash matrix (`CrashRecoveryTest`, see [development.md](development.md)), not by this
benchmark.

**Bytes written.** For every write workload the report also shows the bytes each engine wrote per operation, in a
*Resources* table, plus the total for the whole run. Both engines count the bytes they write themselves:

Expand Down
Loading