From f1ecbaa517f9d0f61aeae551f61490d7e075f658 Mon Sep 17 00:00:00 2001 From: milner Date: Thu, 22 Jun 2017 12:11:00 +0000 Subject: [PATCH] Measure changed-key work on large collections --- Makefile | 9 +- bench/bench.ml | 191 +++++++++++++++++++++++++++++++++++++++ src/emit.ml | 13 ++- test/test_codegen.ml | 6 +- test/test_incremental.ml | 39 ++++++++ test/test_main.ml | 1 + 6 files changed, 252 insertions(+), 7 deletions(-) create mode 100644 bench/bench.ml diff --git a/Makefile b/Makefile index 5f3d97f..d8fd972 100644 --- a/Makefile +++ b/Makefile @@ -149,11 +149,18 @@ $(TEST_EXE): $(TEST_CMX) $(LIB_ARCHIVE) $(RUNTIME_ARCHIVE) $(OCAMLOPT) $(FLAGS) -o $@ $(LIBS) $(RUNTIME_ARCHIVE) $(LIB_ARCHIVE) $(TEST_ORDER) endif +ifeq ($(strip $(BENCH_ML)),) +BENCH_EXE = +else +$(BENCH_EXE): $(BUILD)/bench.cmx $(LIB_ARCHIVE) $(RUNTIME_ARCHIVE) deltac + $(OCAMLOPT) $(FLAGS) -o $@ $(LIBS) $(RUNTIME_ARCHIVE) $(LIB_ARCHIVE) $(BUILD)/bench.cmx +endif + test: $(TEST_EXE) DELTA_ROOT=$(CURDIR) DELTA_BUILD_DIR=$(CURDIR)/$(BUILD) DELTA_OCAMLOPT="$(OCAMLOPT)" ./$(TEST_EXE) bench: $(BENCH_EXE) - DELTA_ROOT=$(CURDIR) DELTA_BUILD_DIR=$(CURDIR)/$(BUILD) ./$(BENCH_EXE) + DELTA_ROOT=$(CURDIR) DELTA_BUILD_DIR=$(CURDIR)/$(BUILD) DELTA_OCAMLOPT="$(OCAMLOPT)" ./$(BENCH_EXE) install: deltac $(RUNTIME_ARCHIVE) $(INSTALL) -d $(DESTDIR)$(BINDIR) $(DESTDIR)$(LIBDIR) diff --git a/bench/bench.ml b/bench/bench.ml new file mode 100644 index 0000000..9ffe6be --- /dev/null +++ b/bench/bench.ml @@ -0,0 +1,191 @@ +let root () = try Sys.getenv "DELTA_ROOT" with Not_found -> "." + +let fail name message = + prerr_endline (Printf.sprintf "bench: %s: %s" name message); + exit 1 + +let deltac () = Filename.concat (root ()) "deltac" + +let build_dir () = try Sys.getenv "DELTA_BUILD_DIR" with Not_found -> "_build" + +let run command = + let log = Filename.temp_file "delta_bench" ".log" in + let fd = Unix.openfile log [ Unix.O_WRONLY; Unix.O_CREAT; Unix.O_TRUNC ] 0o600 in + let argv = Array.of_list command in + let pid = Unix.create_process argv.(0) argv Unix.stdin fd fd in + let status = snd (Unix.waitpid [] pid) in + Unix.close fd; + let text = Native.read_file log in + (try Sys.remove log with _ -> ()); + (status, text) + +let write_file path text = + let channel = open_out path in + output_string channel text; + close_out channel + +type row = { customer : string; total : int } + +let make_row index = { customer = Printf.sprintf "customer%d" index; total = (index * 37) mod 4001 } + +let row_text row = Printf.sprintf "(record (customer %S) (total %d))" row.customer row.total + +let round_of size seed = ((seed * 1103515245) + size) mod size + +let input_text size = + let buffer = Buffer.create (size * 48) in + for index = 0 to size - 1 do + Buffer.add_string buffer (Printf.sprintf "(%d %s)\n" (index + 1) (row_text (make_row index))) + done; + Buffer.contents buffer + +let updates_text size batches = + let buffer = Buffer.create (batches * 64) in + for batch = 0 to batches - 1 do + let key = (round_of size batch) + 1 in + if batch mod 2 = 0 then + Buffer.add_string buffer + (Printf.sprintf "(batch (replace %d %s))\n" key + (row_text { customer = Printf.sprintf "customer%d" key; total = 2000 + (batch mod 500) })) + else + Buffer.add_string buffer + (Printf.sprintf "(batch (insert %d %s))\n" (size + batch + 1) + (row_text { customer = Printf.sprintf "extra%d" batch; total = 1500 + (batch mod 250) })) + done; + Buffer.contents buffer + +type stats = { + st_init_seconds : float; + st_update_seconds : float; + st_minor_words : float; + st_init_counters : string; + st_update_counters : string; +} + +let field prefix line = + let length = String.length prefix in + if String.length line > length && String.sub line 0 length = prefix then + Some (String.trim (String.sub line length (String.length line - length))) + else None + +let counter_value name text = + let parts = String.split_on_char ' ' text in + let rec search = function + | [] -> 0 + | part :: rest -> ( + match (try Some (String.index part '=') with Not_found -> None) with + | Some index + when String.sub part 0 index = name -> + int_of_string (String.sub part (index + 1) (String.length part - index - 1)) + | _ -> search rest) + in + search parts + +let parse_stats text = + let lines = String.split_on_char '\n' text in + let value prefix default = + let rec search = function + | [] -> default + | line :: rest -> ( + match field prefix line with + | Some text -> float_of_string text + | None -> search rest) + in + search lines + in + let counter prefix default = + let rec search = function + | [] -> default + | line :: rest -> ( + match field prefix line with Some text -> text | None -> search rest) + in + search lines + in + { + st_init_seconds = value "init_seconds: " 0.0; + st_update_seconds = value "update_seconds: " 0.0; + st_minor_words = value "minor_words: " 0.0; + st_init_counters = counter "init_counters: " ""; + st_update_counters = counter "update_counters: " ""; + } + +let query_source name = + let path = Filename.concat (Filename.concat (root ()) "example") (name ^ ".delta") in + Native.read_file path + +let compile_query name directory = + let source = query_source name in + let program = Filename.concat directory (name ^ ".ml") in + write_file program source; + let plan = Simplify.simplify (Graph.build (Anf.program (Specialize.program (Infer.program (Resolve.program (Parse.program source)))))) in + let emitted = Filename.concat directory (name ^ "_generated.ml") in + write_file emitted (Emit.program_to_string plan); + let executable = Filename.concat directory name in + let command = + [ (try Sys.getenv "DELTA_OCAMLOPT" with Not_found -> "ocamlopt"); + "-w"; "-26"; "-I"; build_dir (); "-I"; "+unix"; "-o"; executable; + Filename.concat (build_dir ()) "delta_runtime.cmxa"; "unix.cmxa"; emitted ] + in + match run command with + | Unix.WEXITED 0, _ -> executable + | _, output -> fail "compile" output + +let sizes = [ 1000; 100000; 1000000 ] + +let queries = [ "expensive_order"; "revenue"; "count_large" ] + +let () = + let directory = Filename.temp_file "delta_bench" "" in + Sys.remove directory; + Unix.mkdir directory 0o700; + print_endline "delta benchmarks"; + print_endline ""; + print_endline + "query rows init_ms update_ms batches update_us_per_batch minor_words predicate_evals mapping_evals changed_keys full_traversals"; + let regression_failures = ref 0 in + List.iter + (fun name -> + let executable = compile_query name directory in + List.iter + (fun size -> + let input_path = Filename.concat directory (Printf.sprintf "input_%s_%d.sexp" name size) in + let updates_path = Filename.concat directory (Printf.sprintf "update_%s_%d.sexp" name size) in + write_file input_path (input_text size); + write_file updates_path (updates_text size 200); + let status, output = + run [ executable; "--input"; input_path; "--updates"; updates_path; "--print-result"; "--stats" ] + in + if status <> Unix.WEXITED 0 then fail "benchmark run" output; + let stats = parse_stats output in + let batches = 200 in + Printf.printf "%-17s %-9d %-8.3f %-9.3f %-8d %-20.3f %-12.0f %-16d %-14d %-13d %d\n" name size + (stats.st_init_seconds *. 1000.0) (stats.st_update_seconds *. 1000.0) batches + (stats.st_update_seconds *. 1000000.0 /. float_of_int batches) + stats.st_minor_words + (counter_value "predicate_evaluations" stats.st_update_counters) + (counter_value "mapping_evaluations" stats.st_update_counters) + (counter_value "changed_key_visits" stats.st_update_counters) + (counter_value "full_traversals" stats.st_update_counters); + if counter_value "full_traversals" stats.st_update_counters <> 0 then ( + incr regression_failures; + Printf.printf " regression: updates performed %d full traversals for %s at %d rows\n" + (counter_value "full_traversals" stats.st_update_counters) name size); + if counter_value "predicate_evaluations" stats.st_update_counters > 200 * 4 then ( + incr regression_failures; + Printf.printf " regression: too many predicate evaluations (%d) for %s at %d rows\n" + (counter_value "predicate_evaluations" stats.st_update_counters) name size); + if counter_value "mapping_evaluations" stats.st_update_counters > 200 * 4 then ( + incr regression_failures; + Printf.printf " regression: too many mapping evaluations (%d) for %s at %d rows\n" + (counter_value "mapping_evaluations" stats.st_update_counters) name size); + Sys.remove input_path; + Sys.remove updates_path) + sizes) + queries; + print_endline ""; + if !regression_failures = 0 then + print_endline + "regression check: single-key updates performed no full traversals and touched only the changed keys" + else Printf.printf "%d regression checks failed\n" !regression_failures; + Native.remove_dir directory; + exit (if !regression_failures = 0 then 0 else 1) diff --git a/src/emit.ml b/src/emit.ml index fde27d6..8800f25 100644 --- a/src/emit.ml +++ b/src/emit.ml @@ -908,6 +908,7 @@ let emit_driver plan buffer = line buffer " let print_result = ref false in"; line buffer " let stats = ref false in"; line buffer " let trace = ref false in"; + line buffer " let verify = ref false in"; line buffer " let fail message = prerr_endline (\"error: \" ^ message); exit 1 in"; line buffer " let usage message = prerr_endline (\"usage: \" ^ message); exit 2 in"; line buffer " let rec parse_arguments arguments ="; @@ -918,6 +919,7 @@ let emit_driver plan buffer = line buffer " | \"--print-result\" :: rest -> print_result := true; parse_arguments rest"; line buffer " | \"--stats\" :: rest -> stats := true; parse_arguments rest"; line buffer " | \"--trace\" :: rest -> trace := true; parse_arguments rest"; + line buffer " | \"--verify\" :: rest -> verify := true; parse_arguments rest"; line buffer " | flag :: _ -> usage (\"unknown option \" ^ flag)"; line buffer " in"; line buffer " parse_arguments (List.tl (Array.to_list Sys.argv));"; @@ -933,7 +935,7 @@ let emit_driver plan buffer = line buffer " | Delta_runtime.Failure message -> fail (file ^ \": \" ^ message)"; line buffer (Printf.sprintf - " | Delta_runtime.Success entries ->\n let rec decode acc pending =\n match pending with\n | [] -> List.rev acc\n | (key, value) :: rest ->\n (match Query.decode_row value with\n | Delta_runtime.Failure message ->\n fail (Printf.sprintf \"%%s: key %%d: %%s\" file key message)\n | Delta_runtime.Success row ->\n if Delta_runtime.Pure_map.mem key (List.fold_left (fun map (key, _) -> Delta_runtime.Pure_map.add key () map) Delta_runtime.Pure_map.empty acc) then\n fail (Printf.sprintf \"%%s: duplicate key %%d\" file key)\n else decode ((key, row) :: acc) rest)\n in\n decode [] entries));"); + " | Delta_runtime.Success entries ->\n let rec decode seen acc pending =\n match pending with\n | [] -> List.rev acc\n | (key, value) :: rest ->\n if Delta_runtime.Pure_map.mem key seen then\n fail (Printf.sprintf \"%%s: duplicate key %%d\" file key)\n else\n (match Query.decode_row value with\n | Delta_runtime.Failure message ->\n fail (Printf.sprintf \"%%s: key %%d: %%s\" file key message)\n | Delta_runtime.Success row ->\n decode (Delta_runtime.Pure_map.add key () seen) ((key, row) :: acc) rest)\n in\n decode Delta_runtime.Pure_map.empty [] entries));"); line buffer " in"; line buffer (Printf.sprintf " let state = ref (Query.init (rows : (int * %s) list)) in" input); @@ -970,13 +972,14 @@ let emit_driver plan buffer = line buffer " fail (Printf.sprintf \"%s: batch %d: %s\" file (index + 1) message)"; line buffer " | Delta_runtime.Success (next, change) ->"; + line buffer " if !verify then ("; line buffer - " let applied = Query.apply_output_change (Query.result !state) change in"; - line buffer " let current = Query.result next in"; + " let applied = Query.apply_output_change (Query.result !state) change in"; + line buffer " let current = Query.result next in"; line buffer - " if not (Query.equal_output applied current) then"; + " if not (Query.equal_output applied current) then"; line buffer - " fail (Printf.sprintf \"output change does not reproduce the result after batch %d\" (index + 1));"; + " fail (Printf.sprintf \"output change does not reproduce the result after batch %d\" (index + 1)));"; line buffer " state := next;"; line buffer " update_seconds := !update_seconds +. (Sys.time () -. before);"; line buffer " if !trace then"; diff --git a/test/test_codegen.ml b/test/test_codegen.ml index 470cb91..fe9aba7 100644 --- a/test/test_codegen.ml +++ b/test/test_codegen.ml @@ -659,7 +659,11 @@ let generated_differential_case seed_count batch_count = let updates_path = Filename.concat dir (Printf.sprintf "updates_%d.sexp" seed) in write_entries input_path entries; write_batches updates_path batches; - (match run [ executable; "--input"; input_path; "--updates"; updates_path; "--trace"; "--print-result" ] with + (match run + [ + executable; "--input"; input_path; "--updates"; updates_path; "--trace"; + "--verify"; "--print-result"; + ] with | Unix.WEXITED 0, output -> let lines = List.filter (fun line -> line <> "") (String.split_on_char '\n' output) in if List.length lines <> List.length expected then diff --git a/test/test_incremental.ml b/test/test_incremental.ml index 47b5f20..f13541e 100644 --- a/test/test_incremental.ml +++ b/test/test_incremental.ml @@ -856,3 +856,42 @@ let differential_cases = ( "the incremental plan matches full evaluation on a short run", fun () -> differential_case 5 20 ); ] + +let regression_cases = + [ + ( "a single key update on a large collection touches only that key", + fun () -> + let entries = + Util.list_init 500 (fun index -> (index + 1, order_row (Printf.sprintf "customer%d" index) ((index * 37) mod 4001))) + in + let fixture = build_fixture (read_fixture "expensive_order.delta") entries in + Delta_runtime.reset_counters fixture.fx_counters; + (match step fixture [ Change.OpReplace (250, order_row "customer250" 3000) ] with + | Error message -> fail "update" message + | Ok (applied, cached, _) -> + let counters = fixture.fx_counters in + check_equal_int "predicate evaluations" 1 counters.Delta_runtime.predicate_evaluations; + check_equal_int "mapping evaluations" 0 counters.Delta_runtime.mapping_evaluations; + check_equal_int "scalar deltas" 1 counters.Delta_runtime.scalar_deltas; + check_equal_int "changed key visits" 3 counters.Delta_runtime.changed_key_visits; + check_equal_int "full traversals" 0 counters.Delta_runtime.full_traversals; + check_equal_string "result matches the reference" + (Value.to_string (reference_result fixture fixture.fx_state)) + (Value.to_string applied); + check_equal_string "cache matches the reference" + (Value.to_string (reference_result fixture fixture.fx_state)) + (Value.to_string cached)) ); + ( "an inserted key only evaluates the predicate and mapping for itself", + fun () -> + let entries = Util.list_init 500 (fun index -> (index + 1, order_row "row" 10)) in + let fixture = build_fixture (read_fixture "expensive_order.delta") entries in + Delta_runtime.reset_counters fixture.fx_counters; + (match step fixture [ Change.OpInsert (501, order_row "newcomer" 5000) ] with + | Error message -> fail "insert" message + | Ok _ -> + let counters = fixture.fx_counters in + check_equal_int "predicate evaluations" 1 counters.Delta_runtime.predicate_evaluations; + check_equal_int "mapping evaluations" 1 counters.Delta_runtime.mapping_evaluations; + check_equal_int "changed key visits" 3 counters.Delta_runtime.changed_key_visits; + check_equal_int "full traversals" 0 counters.Delta_runtime.full_traversals) ); + ] diff --git a/test/test_main.ml b/test/test_main.ml index 0dbf493..1cfd639 100644 --- a/test/test_main.ml +++ b/test/test_main.ml @@ -388,6 +388,7 @@ let () = Test_harness.run_suite "generated updates" Test_codegen.update_cases; Test_harness.run_suite "wire" Test_codegen.wire_cases; Test_harness.run_suite "cli" Test_codegen.cli_cases; + Test_harness.run_suite "regression" Test_incremental.regression_cases; Test_harness.run_suite "differential" Test_incremental.differential_cases; Test_harness.run_suite "generated differential" Test_codegen.generated_differential_cases; Printf.printf "%d cases, %d failures\n" (Test_harness.case_count ()) (Test_harness.failure_count ());