From eb225637bfa2b6cc7bd19949200fbfda1a3117a7 Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 16:28:50 +0700 Subject: [PATCH 1/3] exec: queue a batch of jobs under one lock and one broadcast Timing each worker in place on a quiet 32 core box showed a query asking for eight hands regularly getting seven. The last job of the batch would start after the other seven had finished the scan between them, and the six before it started in a ramp, each about a wakeup behind the one before. The pool pushed one job at a time, taking the queue lock and waking one worker for each. A wakeup is not addressed to a job, though. A worker that is already awake takes the job the notify was meant for, the notified worker finds the queue empty and parks again, and the job it should have taken sits there until somebody finishes and loops around. For a morsel scan that is the whole query, since a worker does not come back until the morsels run out. The batch now goes under one lock and one broadcast, so every parked worker wakes with jobs already queued and takes one. Spawning moved outside the lock while it was here: a spawn is long enough that holding the queue through it stalls the workers waiting to dequeue, which is the opposite of what spawning is for. The submitting thread also spins on the latch before it parks now. It is not a worker, so the reason workers park at once does not apply to it: it has just finished its own share and the hands it waits on are inside microseconds of finishing theirs, so the core it holds is one nobody wants. Parking there cost more than a tenth of a millisecond against queries that run in about one. --- crates/zu-exec/src/pool.rs | 134 +++++++++++++++++++++++++++++-------- 1 file changed, 107 insertions(+), 27 deletions(-) diff --git a/crates/zu-exec/src/pool.rs b/crates/zu-exec/src/pool.rs index 45dd7eb7..3d108ccb 100644 --- a/crates/zu-exec/src/pool.rs +++ b/crates/zu-exec/src/pool.rs @@ -33,6 +33,23 @@ const SPIN: std::time::Duration = std::time::Duration::from_millis(1); #[cfg(not(windows))] const SPIN: std::time::Duration = std::time::Duration::ZERO; +/// How long the submitting thread spins on the latch before parking. +/// +/// It is not a worker and the argument above does not apply to it: it +/// has just finished its own share and the hands it is waiting on are +/// inside microseconds of finishing theirs, so the core it holds is +/// one nobody else wants for that stretch. Parking instead costs a +/// wakeup, which timed in place on a quiet 32 core box came to more +/// than a tenth of a millisecond against a query that runs in one. +#[cfg(not(windows))] +const LATCH_SPIN: std::time::Duration = std::time::Duration::from_micros(200); + +/// Windows already spins a whole millisecond everywhere else in here, +/// for the same reason and at four times the cost, so the latch keeps +/// the number it has always used. +#[cfg(windows)] +const LATCH_SPIN: std::time::Duration = SPIN; + struct Queue { jobs: Mutex>, ready: Condvar, @@ -133,7 +150,7 @@ impl Latch { impl Pending { pub(crate) fn wait(&self) { let start = std::time::Instant::now(); - while start.elapsed() < SPIN { + while start.elapsed() < LATCH_SPIN { if self.latch.left.load(Ordering::Acquire) == 0 { return; } @@ -156,6 +173,16 @@ impl Drop for Pending { /// borrows. Each job runs exactly once, panics included: a panicking /// job counts down on unwind so the latch cannot hang, and the caller /// sees the panic as its result slot staying empty. +/// +/// The whole batch goes under one lock and one broadcast. Pushing them +/// one at a time and waking one worker each was a wakeup per job, and a +/// wakeup is not addressed to a job: a worker already awake takes the +/// job the notify was for, the notified worker finds the queue empty +/// and parks again, and the job it should have taken waits for whoever +/// finishes first. Timed per worker on a quiet 32 core box, that left +/// the last of eight starting after the other seven had finished the +/// scan between them, so a query asking for eight hands regularly got +/// seven, and the six before it started in a ramp a wakeup apart. pub(crate) fn submit<'a>(jobs: Vec>) -> Pending { let q = queue(); let latch = Arc::new(Latch { @@ -163,43 +190,53 @@ pub(crate) fn submit<'a>(jobs: Vec>) -> Pending { mu: Mutex::new(()), done: Condvar::new(), }); + if jobs.is_empty() { + return Pending { latch }; + } let cap = std::thread::available_parallelism().map_or(1, |n| n.get()); + let count = jobs.len(); // SAFETY: Pending waits, in wait() or in drop, until every job has // run, so the 'a borrows inside the jobs are live for as long as // any worker can touch them. let jobs: Vec = unsafe { std::mem::transmute(jobs) }; - for job in jobs { - let latch = Arc::clone(&latch); - let wrapped: Job = Box::new(move || { - struct CountDown(Arc); - impl Drop for CountDown { - fn drop(&mut self) { - self.0.count_down(); + let spawn = { + let mut queued = q.jobs.lock().unwrap(); + for job in jobs { + let latch = Arc::clone(&latch); + let wrapped: Job = Box::new(move || { + struct CountDown(Arc); + impl Drop for CountDown { + fn drop(&mut self) { + self.0.count_down(); + } } - } - let _count = CountDown(latch); - job(); - }); - let mut jobs = q.jobs.lock().unwrap(); - jobs.push_back(wrapped); - q.pending.fetch_add(1, Ordering::Relaxed); + let _count = CountDown(latch); + job(); + }); + queued.push_back(wrapped); + } + let pending = q.pending.fetch_add(count, Ordering::Relaxed) + count; // Spawn while queued jobs outnumber the workers free to take // them; comparing against zero here left one parked worker // serving a whole batch alone, because idle only drops when a - // worker actually dequeues, long after this loop has pushed + // worker actually dequeues, long after this has queued // everything. - if q.pending.load(Ordering::Relaxed) > q.idle.load(Ordering::Relaxed) - && q.workers.load(Ordering::Relaxed) < cap - { - q.workers.fetch_add(1, Ordering::Relaxed); - std::thread::Builder::new() - .name("zu-exec-worker".into()) - .spawn(move || worker_loop(q)) - .expect("spawn a pool worker"); - } - drop(jobs); - q.ready.notify_one(); + let free = q.idle.load(Ordering::Relaxed); + let room = cap.saturating_sub(q.workers.load(Ordering::Relaxed)); + let spawn = pending.saturating_sub(free).min(room); + q.workers.fetch_add(spawn, Ordering::Relaxed); + spawn + }; + // Outside the lock: a spawn takes long enough that holding the + // queue through it would stall the workers already waiting to + // dequeue, which is the opposite of what spawning is for. + for _ in 0..spawn { + std::thread::Builder::new() + .name("zu-exec-worker".into()) + .spawn(move || worker_loop(q)) + .expect("spawn a pool worker"); } + q.ready.notify_all(); Pending { latch } } @@ -224,6 +261,49 @@ mod tests { } } + /// Every job in a batch has to get a worker of its own, which is + /// what a query asking for eight hands is asking for. Each job here + /// counts itself in, waits for the rest to arrive, and counts itself + /// out, so the high water mark is how many of them the pool had + /// running at once. A batch that leaves one job in the queue while + /// another worker runs two of them never reaches the mark, and the + /// wait gives up on a timeout rather than hanging. + #[test] + fn a_batch_gets_a_worker_apiece() { + let hands = std::thread::available_parallelism() + .map_or(2, |n| n.get()) + .min(4); + // Running now, the high water mark, and whether the mark was + // ever reached, which the wait has to watch instead of the + // count: the last job to arrive is also the first to leave. + let state = (Mutex::new((0usize, 0usize, false)), Condvar::new()); + let jobs: Vec> = (0..hands) + .map(|_| { + let state = &state; + Box::new(move || { + let (mu, all_here) = state; + let mut count = mu.lock().unwrap(); + count.0 += 1; + count.1 = count.1.max(count.0); + count.2 |= count.0 == hands; + all_here.notify_all(); + while !count.2 { + let (next, gave_up) = all_here + .wait_timeout(count, std::time::Duration::from_secs(5)) + .unwrap(); + count = next; + if gave_up.timed_out() { + break; + } + } + count.0 -= 1; + }) as Box + }) + .collect(); + submit(jobs).wait(); + assert_eq!(state.0.lock().unwrap().1, hands, "all {hands} ran at once"); + } + #[test] fn a_panicking_job_still_counts_down() { let ok = Mutex::new(false); From acea1ad1def0e260e8addecd238f45f6eb8b8916 Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 16:32:55 +0700 Subject: [PATCH 2/3] exec: let an idle worker spin before it parks on Linux too The spin was a Windows-only thing, on the reading that a parked thread wakes in microseconds everywhere else. Timing every worker in place said otherwise. On a quiet 32 core Linux box a batch of eight started in a ramp of about twenty five microseconds a worker, and the last one or two of them regularly did not arrive at all until the others had finished the scan between them, which is a query asking for eight hands and getting six. A tenth of a millisecond is not microseconds against a query that runs in one. With a hundred microseconds of spin every one of the eight is inside the scan and starts within a tenth of a millisecond of the first, and the stragglers are gone. A tenth of what Windows spins, because the wakeups here cost that much less and because a spinning worker is a worker not parked, which matters on a host with fewer cores than the query asked for. --- crates/zu-exec/src/pool.rs | 20 +++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/crates/zu-exec/src/pool.rs b/crates/zu-exec/src/pool.rs index 3d108ccb..1ab9a2cd 100644 --- a/crates/zu-exec/src/pool.rs +++ b/crates/zu-exec/src/pool.rs @@ -28,10 +28,24 @@ type Job = Box; #[cfg(windows)] const SPIN: std::time::Duration = std::time::Duration::from_millis(1); -/// Elsewhere a parked thread wakes in microseconds and spinning only -/// steals cores from the query still running, so workers park at once. +/// Elsewhere it used to be zero, on the reading that a parked thread +/// wakes in microseconds and that spinning only steals cores from the +/// query still running. Timing every worker in place said otherwise: +/// on a quiet 32 core Linux box, a batch of eight started in a ramp of +/// about twenty five microseconds a worker and the last one or two of +/// them regularly did not arrive until the others had finished the scan +/// between them. A tenth of a millisecond is not microseconds against a +/// query that runs in one. +/// +/// A tenth of what Windows takes, because the wakeups here are that +/// much cheaper and because a worker that spins is a worker not parked +/// on a host with fewer cores than the query wanted. Long enough to +/// cover the gap between one query and the next in a session, which is +/// the case that matters: a worker that stayed awake is the difference +/// between a query getting the hands it asked for and getting most of +/// them. #[cfg(not(windows))] -const SPIN: std::time::Duration = std::time::Duration::ZERO; +const SPIN: std::time::Duration = std::time::Duration::from_micros(100); /// How long the submitting thread spins on the latch before parking. /// From 1f32c33bb633f0faa8c20dcf65d8b446da8539c0 Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 16:37:15 +0700 Subject: [PATCH 3/3] bench: write down what the pool was costing a parallel query --- bench/budgets.toml | 24 ++++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/bench/budgets.toml b/bench/budgets.toml index cdc83299..577b9e1c 100644 --- a/bench/budgets.toml +++ b/bench/budgets.toml @@ -681,6 +681,30 @@ vec_unpack_gvals_s = 0.9 # 3.9x with the 1-worker time itself moving 12 percent run to run. # Nothing in the fleet is 8 quiet equal cores, so the floor sits under # server3's worst run. Raise this floor, never lower it. +# +# Timing every worker in place rather than the query around them found +# part of what those numbers were measuring. The pool queued a batch a +# job at a time and woke one worker for each, and a wakeup is not +# addressed to a job: a worker already awake takes the job the notify +# was meant for, the notified worker finds the queue empty and parks +# again, and the job it should have taken waits for whoever finishes +# first, which for a morsel scan is the end of the query. On gamingpc +# under WSL the last of eight regularly started after the other seven +# had finished the scan between them, and the six before it started in +# a ramp about twenty five microseconds apart, because a parked thread +# there does not wake in microseconds the way this pool assumed. The +# batch goes under one lock and one broadcast now, and an idle worker +# spins a hundred microseconds before it parks on Linux as it already +# did on Windows. All eight are inside the scan after that, starting +# within a tenth of a millisecond of the first. +# +# gamingpc under WSL, medians of nine paired alternating runs of two +# separately built binaries: the expand query went 4.9x to 5.4x at +# eight workers, 1.12 ms to 1.02, faster in all nine pairs, and the +# scan query beside it went 4.2x to 5.3x, 0.80 ms to 0.67. The one +# worker times moved 2 percent the wrong way on both, which is a +# single worker run not going near the pool at all. The same box +# measured 3.6 to 3.7x when the note above was written. exec_scale_8x = 2.0 # The grouping sink at one worker, M rows/s (perf/05 section 4): ten