feat(db-bench): mix + msgrate — concurrent modes (T3)
- Mixer actors: 90/10 read/write, per-actor histograms merged through the store itself (Hist rows) — exact aggregate percentiles - msgrate: one-way flood at a worker-placed sink; measured 15.3M msgs/s same-heap vs 2.06M cross-shard — the mutex-inbox number - finding: point lookups are O(table) (probe walks all slabs), so read-heavy mix is quadratic in store size — all-mode calibrated to N/10 mix ops; the number 22 exists to publish - TSan clean both shard counts (setarch -R, fibers-gate pattern) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
c35e7219c8
commit
7cd56c9426
2 changed files with 187 additions and 2 deletions
|
|
@ -288,6 +288,152 @@ fn verify_acked(m: Int) -> Int {
|
|||
return 0;
|
||||
}
|
||||
|
||||
-- ---- the concurrent modes (mix, msgrate) ----
|
||||
|
||||
fn hist_dump(h: map<Int, Int>, kind: Int) {
|
||||
let b = 0;
|
||||
while b <= 20000 {
|
||||
if has(h, b) {
|
||||
insert Hist { kind: kind, b: b, c: get(h, b) };
|
||||
}
|
||||
b = b + 1;
|
||||
}
|
||||
}
|
||||
|
||||
-- One mixer = one actor: 90/10 read/write over the seeded store. On a
|
||||
-- worker shard every statement below rides the stage-3 DB RPC — the
|
||||
-- code must not know or care (transparency is the point). Done signal:
|
||||
-- a Meta row main polls for (the coordination idiom this side of 31).
|
||||
class Mixer {
|
||||
id: Int
|
||||
fn receive(msg: MixJob) {
|
||||
let hr: map<Int, Int> = {};
|
||||
let hw: map<Int, Int> = {};
|
||||
let s = msg.seed;
|
||||
let sink = 0;
|
||||
let i = 0;
|
||||
while i < msg.ops {
|
||||
s = lcg(s);
|
||||
let key = s % msg.kmod;
|
||||
let o0 = time.ticks();
|
||||
if i % 10 == 9 {
|
||||
let xs = from x in Item where x.k == key take 1 select x;
|
||||
if len(xs) > 0 {
|
||||
xs[0].v = xs[0].v + 1;
|
||||
}
|
||||
hist_add(hw, time.ticks() - o0);
|
||||
} else {
|
||||
let xs = from x in Item where x.k == key take 1 select x;
|
||||
if len(xs) > 0 {
|
||||
sink = sink + xs[0].v;
|
||||
}
|
||||
hist_add(hr, time.ticks() - o0);
|
||||
}
|
||||
i = i + 1;
|
||||
}
|
||||
hist_dump(hr, 0);
|
||||
hist_dump(hw, 1);
|
||||
insert Meta { tag: "mixdone${self.id}", val: sink };
|
||||
}
|
||||
}
|
||||
|
||||
fn mix_mode(total: Int, c: Int) -> Int {
|
||||
let kmod = meta_val("kmod");
|
||||
if kmod < 1 {
|
||||
print_err("mix: seed first");
|
||||
return 1;
|
||||
}
|
||||
let per = total / c;
|
||||
if per < 1 {
|
||||
per = 1;
|
||||
}
|
||||
let wall0 = time.ticks();
|
||||
let i = 0;
|
||||
while i < c {
|
||||
let a: actor MixJob = spawn Mixer { id: i };
|
||||
send(a, MixJob { ops: per, seed: 1000 + i * 7919, kmod: kmod });
|
||||
i = i + 1;
|
||||
}
|
||||
-- poll until every mixer's done row exists
|
||||
let done = 0;
|
||||
while done < c {
|
||||
time.sleep(20);
|
||||
done = 0;
|
||||
i = 0;
|
||||
while i < c {
|
||||
if meta_val("mixdone${i}") >= 0 {
|
||||
done = done + 1;
|
||||
}
|
||||
i = i + 1;
|
||||
}
|
||||
}
|
||||
let wall = time.ticks() - wall0;
|
||||
-- merge the dumped histograms; wall time is shared by both classes
|
||||
let hr: map<Int, Int> = {};
|
||||
let hw: map<Int, Int> = {};
|
||||
let nr = 0;
|
||||
let nw = 0;
|
||||
for x in from x in Hist select x {
|
||||
if x.kind == 0 {
|
||||
if has(hr, x.b) {
|
||||
set(hr, x.b, get(hr, x.b) + x.c);
|
||||
} else {
|
||||
set(hr, x.b, x.c);
|
||||
}
|
||||
nr = nr + x.c;
|
||||
} else {
|
||||
if has(hw, x.b) {
|
||||
set(hw, x.b, get(hw, x.b) + x.c);
|
||||
} else {
|
||||
set(hw, x.b, x.c);
|
||||
}
|
||||
nw = nw + x.c;
|
||||
}
|
||||
}
|
||||
report("mixread", nr, wall, hr);
|
||||
report("mixwrite", nw, wall, hw);
|
||||
return 0;
|
||||
}
|
||||
|
||||
-- msgrate: one-way flood — main sends N messages at a sink actor; the
|
||||
-- sink counts and writes the done row at N. Spawn TWO sinks and flood
|
||||
-- the second: round-robin placement puts it off the primary whenever
|
||||
-- more than one shard exists, so the multi-shard number prices the
|
||||
-- mutex inbox (stage-2 deviation 4's number); single-shard prices the
|
||||
-- same-heap path.
|
||||
class Sink {
|
||||
got: Int
|
||||
fn receive(msg: Flood) {
|
||||
self.got = self.got + 1;
|
||||
if self.got == msg.n {
|
||||
insert Meta { tag: "flooddone", val: self.got };
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn msgrate_mode(n: Int) -> Int {
|
||||
let first: actor Flood = spawn Sink { got: 0 };
|
||||
let a: actor Flood = spawn Sink { got: 0 };
|
||||
if first == a {
|
||||
print_err("msgrate: impossible");
|
||||
}
|
||||
let t0 = time.ticks();
|
||||
let i = 0;
|
||||
while i < n {
|
||||
send(a, Flood { n: n });
|
||||
i = i + 1;
|
||||
}
|
||||
while meta_val("flooddone") < 0 {
|
||||
time.sleep(5);
|
||||
}
|
||||
let us = time.ticks() - t0;
|
||||
if us < 1 {
|
||||
us = 1;
|
||||
}
|
||||
print("msgrate ${n} ${n * 1000000 / us}");
|
||||
return 0;
|
||||
}
|
||||
|
||||
-- all N: the throughput campaign in ONE process — without WO_DATA the
|
||||
-- store is RAM and dies with the process, so seed and the measured
|
||||
-- modes must share a run; under WO_DATA the same mode prices the
|
||||
|
|
@ -309,13 +455,17 @@ fn all_mode(n: Int) -> Int {
|
|||
if rc != 0 {
|
||||
return rc;
|
||||
}
|
||||
return 0;
|
||||
-- mix at N/10: every point lookup is O(table) today (the probe walks
|
||||
-- all slabs — a headline finding, not a bug to hide), so a read-heavy
|
||||
-- mix over a seeded store is quadratic in N. The campaign driver
|
||||
-- chooses absolute sizes; this keeps `all` finishing in minutes.
|
||||
return mix_mode(n / 10, 4);
|
||||
}
|
||||
|
||||
fn usage() -> Int {
|
||||
print_err("usage: db-bench <mode>");
|
||||
print_err(" all N | seed N | read N | query N | write N | wal N");
|
||||
print_err(" verify | verify-acked M");
|
||||
print_err(" mix N C | msgrate N | verify | verify-acked M");
|
||||
return 2;
|
||||
}
|
||||
|
||||
|
|
@ -355,5 +505,19 @@ fn main(args: multi Text) -> Int {
|
|||
if args[0] == "verify-acked" {
|
||||
return verify_acked(n);
|
||||
}
|
||||
if args[0] == "msgrate" {
|
||||
return msgrate_mode(n);
|
||||
}
|
||||
if args[0] == "mix" {
|
||||
if len(args) < 3 {
|
||||
return usage();
|
||||
}
|
||||
let c = parse_int(args[2]);
|
||||
if c == nil or c < 1 {
|
||||
print_err("db-bench: <c> must be a positive number");
|
||||
return 2;
|
||||
}
|
||||
return mix_mode(n, c);
|
||||
}
|
||||
return usage();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,3 +22,24 @@ class Meta {
|
|||
tag: Text @unique
|
||||
val: Int
|
||||
}
|
||||
|
||||
-- mix actors dump their per-op histograms here (kind 0 = read,
|
||||
-- 1 = write); main scans and merges — exact aggregate percentiles,
|
||||
-- and the merge itself dogfoods the store.
|
||||
@table(name: "hist")
|
||||
class Hist {
|
||||
kind: Int
|
||||
b: Int
|
||||
c: Int
|
||||
}
|
||||
|
||||
-- messages (Int-only payloads: ownership moves, nothing borrowed)
|
||||
class MixJob {
|
||||
ops: Int
|
||||
seed: Int
|
||||
kmod: Int
|
||||
}
|
||||
|
||||
class Flood {
|
||||
n: Int
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue