diff --git a/benchmarks/src/main/java/io/hstore/bench/Comparison.java b/benchmarks/src/main/java/io/hstore/bench/Comparison.java index 7998869..07593ad 100644 --- a/benchmarks/src/main/java/io/hstore/bench/Comparison.java +++ b/benchmarks/src/main/java/io/hstore/bench/Comparison.java @@ -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; @@ -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; @@ -76,10 +79,12 @@ public static void main(String[] args) throws Exception { Map 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); } @@ -103,16 +108,11 @@ private static void run(Map 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 results = new ArrayList<>(); Map 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 { @@ -196,11 +196,74 @@ private static void run(Map 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 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 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 options, Dataset dataset) throws Exception { + Path directory = Files.createTempDirectory("hstore-bench-crash"); + List 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)); } diff --git a/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java b/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java index 236d666..66c70f5 100644 --- a/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java +++ b/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java @@ -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 { @@ -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; @@ -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 @@ -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; @@ -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(); diff --git a/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java b/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java index 4fec7a9..5d48958 100644 --- a/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java +++ b/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java @@ -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 { @@ -78,7 +79,7 @@ private T read(Callable 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; @@ -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(); diff --git a/benchmarks/src/main/java/io/hstore/bench/Report.java b/benchmarks/src/main/java/io/hstore/bench/Report.java index 3f2d50d..a81eaab 100644 --- a/benchmarks/src/main/java/io/hstore/bench/Report.java +++ b/benchmarks/src/main/java/io/hstore/bench/Report.java @@ -71,12 +71,13 @@ static String markdown(String[] args) { return out.toString(); } - private record Metric(String key, String name, boolean perOperation, DoubleFunction format) { + private record Metric(String key, String name, boolean perOperation, boolean showZero, DoubleFunction format) { } - private static final List 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 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 subjects, List baselines) { List all = Stream.concat(subjects.stream(), baselines.stream()).toList(); @@ -86,7 +87,7 @@ private static void resources(StringBuilder out, List subjects, List 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; diff --git a/benchmarks/src/main/java/io/hstore/bench/Store.java b/benchmarks/src/main/java/io/hstore/bench/Store.java index 01898aa..8287e8d 100644 --- a/benchmarks/src/main/java/io/hstore/bench/Store.java +++ b/benchmarks/src/main/java/io/hstore/bench/Store.java @@ -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 { @@ -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); diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 237f6d4..87f3e80 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -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`. @@ -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: