The `obikmer select` CLI no longer supports in-place index rewriting; the `--output` flag is now required, with benchmarks updated to use temporary directories for atomic replacement. Added `--dense` and `--force-copy` flags. Introduced `batch_presence_counts` to compute presence counts across multiple column groups in a single pass, eliminating redundant I/O. Refactored the aggregation pipeline to branch on layer content, applying the optimized batched counting for `Presence` layers. Enhanced the NUMA runner to catch worker panics, track them, and re-raise after thread join to prevent indefinite blocking.
453 lines
19 KiB
Rust
453 lines
19 KiB
Rust
use std::any::Any;
|
|
use std::panic::{self, AssertUnwindSafe};
|
|
use std::sync::Arc;
|
|
use std::time::{Duration, Instant};
|
|
|
|
use crossbeam_channel::unbounded;
|
|
use crate::{CpuSample, IoSample};
|
|
use tracing::debug;
|
|
|
|
use super::topology::{build, pin_current_thread};
|
|
|
|
// ── PartitionRunner ─────────────────────────────────────────────────────────
|
|
|
|
/// Growth step (fraction of a node's worker capacity added per activation
|
|
/// event, see [`NodeActivation::grow`]).
|
|
const GROWTH_DIVISOR: usize = 8;
|
|
/// Minimum CPU efficiency growth to activate more workers, as a fraction of
|
|
/// the size of the *last growth step* (e.g. `0.2` after adding 8 workers
|
|
/// requires the next check to show at least +1.6 cores of growth — 20 % of
|
|
/// the ~8 cores those 8 workers should contribute if the workload is truly
|
|
/// CPU-bound). Scaling by the last step's size — not the cumulative total —
|
|
/// keeps the bar meaningful regardless of how many workers are already
|
|
/// active, instead of demanding an ever-larger absolute jump as the pool
|
|
/// grows.
|
|
const CPU_SPAWN_THRESHOLD: f64 = 0.2;
|
|
/// Minimum I/O throughput growth (relative) to activate more workers.
|
|
const IO_SPAWN_THRESHOLD: f64 = 0.2;
|
|
|
|
struct NodeConfig {
|
|
pool: Option<Arc<rayon::ThreadPool>>,
|
|
cpu_ids: Vec<usize>,
|
|
max_workers: usize,
|
|
}
|
|
|
|
/// Generic NUMA-aware runner for partition-level parallel work.
|
|
///
|
|
/// Workers are distributed evenly across NUMA nodes and pinned to their
|
|
/// node's CPUs. UMA is the degenerate case: one node, no pinning.
|
|
///
|
|
/// Workers are pre-spawned dormant, one activation channel per node so
|
|
/// growth always targets a specific node rather than whichever dormant
|
|
/// worker happens to wake up first on a shared channel. Growth (both the
|
|
/// initial count and each subsequent step) is expressed as a fraction of
|
|
/// `workers_per_node`, applied identically to every node, so the pace of
|
|
/// ramp-up depends on node size rather than node count — a single-NUMA-node
|
|
/// (UMA) machine ramps just as fast as an 8-node one.
|
|
///
|
|
/// # Termination
|
|
///
|
|
/// ```text
|
|
/// drop(part_tx) → part_rx drains → workers exit → drop their result_tx
|
|
/// drop(result_tx) → result_rx closes → controller loop exits
|
|
/// drop(activate_txs) → dormant workers exit cleanly
|
|
/// ```
|
|
pub struct PartitionRunner {
|
|
nodes: Vec<NodeConfig>,
|
|
}
|
|
|
|
impl PartitionRunner {
|
|
/// Total worker slots across all nodes.
|
|
pub fn max_workers(&self) -> usize {
|
|
self.nodes.iter().map(|n| n.max_workers).sum()
|
|
}
|
|
|
|
/// Detect topology and build. Always succeeds.
|
|
pub fn new() -> Self {
|
|
let ns = build();
|
|
let wpn = ns.workers_per_node();
|
|
debug!(
|
|
"PartitionRunner: {} node(s) × {} worker(s)/node max",
|
|
ns.pools.len(),
|
|
wpn,
|
|
);
|
|
let nodes = ns
|
|
.pools
|
|
.into_iter()
|
|
.zip(ns.cpus_per_node)
|
|
.map(|(pool, cpu_ids)| NodeConfig {
|
|
pool,
|
|
cpu_ids,
|
|
max_workers: wpn,
|
|
})
|
|
.collect();
|
|
Self { nodes }
|
|
}
|
|
|
|
/// Like [`new`](Self::new), but caps total worker slots (summed across
|
|
/// nodes) at `max_total_workers` — split evenly across nodes, each
|
|
/// further capped by that node's actual core count. For callers whose
|
|
/// own per-worker closure does further internal parallel work (so the
|
|
/// natural per-node core count would oversubscribe if used as the
|
|
/// *outer* degree of parallelism too).
|
|
pub fn new_capped(max_total_workers: usize) -> Self {
|
|
let ns = build();
|
|
let n_nodes = ns.pools.len().max(1);
|
|
let per_node_cap = (max_total_workers / n_nodes).max(1);
|
|
debug!(
|
|
"PartitionRunner (capped): {} node(s) × up to {} worker(s)/node ({} total requested)",
|
|
n_nodes, per_node_cap, max_total_workers,
|
|
);
|
|
let nodes = ns
|
|
.pools
|
|
.into_iter()
|
|
.zip(ns.cpus_per_node)
|
|
.map(|(pool, cpu_ids)| {
|
|
let node_cores = cpu_ids.len().max(1);
|
|
NodeConfig { pool, cpu_ids, max_workers: per_node_cap.min(node_cores) }
|
|
})
|
|
.collect();
|
|
Self { nodes }
|
|
}
|
|
|
|
/// Run `f(i)` for every index in `order`.
|
|
///
|
|
/// Workers are pre-spawned dormant and activated adaptively, per node:
|
|
/// `(workers_per_node / INITIAL_DIVISOR).max(1)` are woken immediately on
|
|
/// every node, then `(workers_per_node / GROWTH_DIVISOR).max(1)` more per
|
|
/// node each time the check below fires. A timer thread fires that check
|
|
/// every `TIMER_SECS` seconds; each completed partition resets that timer
|
|
/// (forcing an immediate check) and also triggers its own inline check. A
|
|
/// growth step happens whenever CPU efficiency grows by at least
|
|
/// `CPU_SPAWN_THRESHOLD` of what the last growth step should have
|
|
/// contributed, or I/O throughput grows by at least `IO_SPAWN_THRESHOLD`
|
|
/// (relative) since the last check — whichever resource is the actual
|
|
/// bottleneck still shows headroom.
|
|
///
|
|
/// `on_done(i, result, elapsed)` is called from the controller thread as
|
|
/// each partition completes — suitable for progress bars and result
|
|
/// aggregation.
|
|
///
|
|
/// Returns the first error produced by `f`, if any.
|
|
pub fn run<F, R, E, C>(&self, order: &[usize], f: F, mut on_done: C) -> Result<(), E>
|
|
where
|
|
F: Fn(usize) -> Result<R, E> + Send + Sync,
|
|
R: Send,
|
|
E: Send,
|
|
C: FnMut(usize, R, Duration) + Send,
|
|
{
|
|
let n_total = order.len();
|
|
if n_total == 0 {
|
|
return Ok(());
|
|
}
|
|
|
|
const TIMER_SECS: u64 = 30;
|
|
const INITIAL_DIVISOR: usize = 4;
|
|
|
|
// ── Channels ──────────────────────────────────────────────────────────
|
|
let (part_tx, part_rx) = unbounded::<usize>();
|
|
// reset_tx: controller → timer ("reset the 30 s window")
|
|
let (reset_tx, reset_rx) = unbounded::<()>();
|
|
// event_tx: workers + timer → controller (unified event stream)
|
|
let (event_tx, event_rx) = unbounded::<WorkerEvent<R, E>>();
|
|
// One activation channel per node: growth always targets a specific
|
|
// node, rather than whichever dormant worker happens to win the race
|
|
// on a channel shared across all nodes.
|
|
let (activate_txs, activate_rxs): (Vec<_>, Vec<_>) =
|
|
(0..self.nodes.len()).map(|_| unbounded::<()>()).unzip();
|
|
|
|
for &i in order {
|
|
part_tx.send(i).ok();
|
|
}
|
|
drop(part_tx);
|
|
|
|
let max_workers = self.max_workers();
|
|
let node_caps: Vec<usize> = self.nodes.iter().map(|n| n.max_workers).collect();
|
|
let f = &f;
|
|
|
|
let mut first_err: Option<E> = None;
|
|
let mut first_panic: Option<Box<dyn Any + Send + 'static>> = None;
|
|
|
|
std::thread::scope(|s| {
|
|
// ── Timer thread ──────────────────────────────────────────────────
|
|
// Sends TimerTick every TIMER_SECS seconds. Resets its window each
|
|
// time reset_rx receives a message (i.e. on partition completion).
|
|
let timer_tx = event_tx.clone();
|
|
s.spawn(move || {
|
|
let period = Duration::from_secs(TIMER_SECS);
|
|
loop {
|
|
crossbeam_channel::select! {
|
|
recv(reset_rx) -> r => {
|
|
if r.is_err() { break; } // reset_tx dropped → exit
|
|
}
|
|
default(period) => {
|
|
if timer_tx.send(WorkerEvent::TimerTick).is_err() { break; }
|
|
}
|
|
}
|
|
}
|
|
});
|
|
|
|
// ── Pre-spawn workers dormant, grouped by node ────────────────────
|
|
// Each worker listens on its own node's activation channel only.
|
|
for (node, arx) in self.nodes.iter().zip(activate_rxs.iter()) {
|
|
let cpu_ids = &node.cpu_ids;
|
|
for _ in 0..node.max_workers {
|
|
let prx = part_rx.clone();
|
|
let etx = event_tx.clone();
|
|
let arx = arx.clone();
|
|
let pool = node.pool.clone();
|
|
|
|
s.spawn(move || {
|
|
let tid = std::thread::current().id();
|
|
debug!(?tid, "PartitionRunner worker: waiting on activation");
|
|
if arx.recv().is_err() {
|
|
debug!(?tid, "PartitionRunner worker: activation channel closed, exiting");
|
|
return;
|
|
}
|
|
debug!(?tid, "PartitionRunner worker: activated");
|
|
if !cpu_ids.is_empty() {
|
|
pin_current_thread(cpu_ids);
|
|
}
|
|
for i in &prx {
|
|
debug!(?tid, partition = i, "PartitionRunner worker: picked partition");
|
|
let t = Instant::now();
|
|
// Caught, not left to unwind straight through this
|
|
// spawned thread: a panicking `f(i)` would
|
|
// otherwise never reach `etx.send(...)` below, so
|
|
// the controller's `completed < n_total` loop
|
|
// waits forever for an event this partition can
|
|
// no longer produce (see `run`'s own docs on the
|
|
// termination protocol) — silently hanging
|
|
// instead of surfacing the panic. Caught here and
|
|
// re-raised on the caller's thread once `run`
|
|
// returns, so the original message/backtrace
|
|
// still surfaces, just from the right place.
|
|
let outcome = panic::catch_unwind(AssertUnwindSafe(|| match &pool {
|
|
Some(p) => p.install(|| f(i)),
|
|
None => f(i),
|
|
}));
|
|
debug!(?tid, partition = i, "PartitionRunner worker: partition done");
|
|
let event = match outcome {
|
|
Ok(r) => WorkerEvent::Completed(i, r, t.elapsed()),
|
|
Err(payload) => WorkerEvent::Panicked(i, payload),
|
|
};
|
|
etx.send(event).ok();
|
|
}
|
|
debug!(?tid, "PartitionRunner worker: no more partitions, exiting");
|
|
});
|
|
}
|
|
}
|
|
// Drop controller's event_tx: event_rx closes when all workers +
|
|
// timer have exited.
|
|
drop(event_tx);
|
|
|
|
// ── Controller ────────────────────────────────────────────────────
|
|
let mut activation = NodeActivation::new(&activate_txs, &node_caps, max_workers);
|
|
activation.activate_initial(INITIAL_DIVISOR, n_total);
|
|
debug!(n_total, activated = activation.total(), "PartitionRunner controller: initial activation");
|
|
|
|
let mut cpu_sample = CpuSample::now();
|
|
let mut io_sample = IoSample::now();
|
|
let mut completed = 0usize;
|
|
|
|
while completed < n_total {
|
|
debug!(completed, n_total, "PartitionRunner controller: waiting for an event");
|
|
let Ok(event) = event_rx.recv() else {
|
|
debug!("PartitionRunner controller: event channel closed, stopping");
|
|
break;
|
|
};
|
|
match event {
|
|
WorkerEvent::Completed(i, r, dur) => {
|
|
match r {
|
|
Ok(v) => on_done(i, v, dur),
|
|
Err(e) => {
|
|
if first_err.is_none() {
|
|
first_err = Some(e);
|
|
}
|
|
}
|
|
}
|
|
completed += 1;
|
|
// Reset the 30 s timer.
|
|
reset_tx.send(()).ok();
|
|
// Inline check: same logic as a timer tick.
|
|
maybe_activate(
|
|
&mut activation,
|
|
&mut cpu_sample,
|
|
&mut io_sample,
|
|
completed,
|
|
n_total,
|
|
);
|
|
}
|
|
WorkerEvent::Panicked(_i, payload) => {
|
|
// Counts toward `completed` like any other outcome —
|
|
// this partition will never produce a `Completed`
|
|
// event, so not counting it here is exactly what
|
|
// used to hang the controller forever. The payload
|
|
// is re-raised on the caller's thread once `run`
|
|
// returns (see below), not here: unwinding out of
|
|
// this `recv` loop would leak the still-running
|
|
// worker/timer threads this `thread::scope` owns.
|
|
if first_panic.is_none() {
|
|
first_panic = Some(payload);
|
|
}
|
|
completed += 1;
|
|
reset_tx.send(()).ok();
|
|
}
|
|
WorkerEvent::TimerTick => {
|
|
maybe_activate(
|
|
&mut activation,
|
|
&mut cpu_sample,
|
|
&mut io_sample,
|
|
completed,
|
|
n_total,
|
|
);
|
|
}
|
|
}
|
|
}
|
|
|
|
// Dormant workers exit once every sender for their node's channel
|
|
// is dropped — `activate_txs` holds the only ones.
|
|
drop(activate_txs);
|
|
// Timer thread exits when reset_tx closes.
|
|
drop(reset_tx);
|
|
});
|
|
|
|
// A panic takes priority over a plain `Err`: it means `f` itself hit
|
|
// a bug (an unhandled case, an assertion) rather than a normal,
|
|
// typed failure — worth surfacing with its original message/
|
|
// backtrace via unwinding, not silently downgraded to whatever `Err`
|
|
// another, unrelated partition happened to return first.
|
|
if let Some(payload) = first_panic {
|
|
panic::resume_unwind(payload);
|
|
}
|
|
|
|
match first_err {
|
|
Some(e) => Err(e),
|
|
None => Ok(()),
|
|
}
|
|
}
|
|
}
|
|
|
|
// ── Internal event type ───────────────────────────────────────────────────────
|
|
|
|
enum WorkerEvent<R, E> {
|
|
Completed(usize, Result<R, E>, Duration),
|
|
/// `f(i)` panicked instead of returning — see `run`'s own docs on why
|
|
/// this is caught at all (never letting the panic unwind straight
|
|
/// through the spawned worker thread) rather than a bug fix that could
|
|
/// be skipped.
|
|
Panicked(usize, Box<dyn Any + Send + 'static>),
|
|
TimerTick,
|
|
}
|
|
|
|
/// Tracks how many of each node's dormant workers have been woken, and
|
|
/// grows every node by the same amount at each step (capped by that node's
|
|
/// remaining dormant workers and by the run's total budget) so load stays
|
|
/// balanced across nodes at every point in time — never just "one more
|
|
/// worker somewhere". Also remembers the size of the last real growth step
|
|
/// (`last_step`), used to scale the CPU activation threshold to what that
|
|
/// step could plausibly have contributed (see `maybe_activate`).
|
|
struct NodeActivation<'a> {
|
|
txs: &'a [crossbeam_channel::Sender<()>],
|
|
caps: &'a [usize],
|
|
active: Vec<usize>,
|
|
total: usize,
|
|
max: usize,
|
|
last_step: usize,
|
|
}
|
|
|
|
impl<'a> NodeActivation<'a> {
|
|
fn new(txs: &'a [crossbeam_channel::Sender<()>], caps: &'a [usize], max: usize) -> Self {
|
|
Self {
|
|
txs,
|
|
caps,
|
|
active: vec![0; txs.len()],
|
|
total: 0,
|
|
max,
|
|
last_step: 0,
|
|
}
|
|
}
|
|
|
|
fn total(&self) -> usize {
|
|
self.total
|
|
}
|
|
fn last_step(&self) -> usize {
|
|
self.last_step
|
|
}
|
|
fn max(&self) -> usize {
|
|
self.max
|
|
}
|
|
fn is_full(&self) -> bool {
|
|
self.total >= self.max
|
|
}
|
|
|
|
/// Wake up to `(node_cap / divisor).max(1)` dormant workers on every
|
|
/// node, capped by `n_total`. Called once at startup, unconditionally.
|
|
fn activate_initial(&mut self, divisor: usize, n_total: usize) {
|
|
self.grow(divisor, n_total);
|
|
}
|
|
|
|
/// Same per-node sizing as [`activate_initial`](Self::activate_initial),
|
|
/// applied as a growth step. Returns the number of workers actually
|
|
/// activated (may be less than requested once a node or the total
|
|
/// budget is exhausted). Updates `last_step` when it actually grew.
|
|
fn grow(&mut self, divisor: usize, n_total: usize) -> usize {
|
|
let before = self.total;
|
|
for idx in 0..self.txs.len() {
|
|
let wanted = (self.caps[idx] / divisor).max(1);
|
|
let room = self.caps[idx].saturating_sub(self.active[idx]);
|
|
let grow = wanted.min(room).min(n_total.saturating_sub(self.total));
|
|
for _ in 0..grow {
|
|
self.txs[idx].send(()).ok();
|
|
}
|
|
self.active[idx] += grow;
|
|
self.total += grow;
|
|
}
|
|
let grew = self.total - before;
|
|
if grew > 0 {
|
|
self.last_step = grew;
|
|
}
|
|
grew
|
|
}
|
|
}
|
|
|
|
fn maybe_activate(
|
|
activation: &mut NodeActivation,
|
|
cpu_sample: &mut CpuSample,
|
|
io_sample: &mut IoSample,
|
|
completed: usize,
|
|
n_total: usize,
|
|
) {
|
|
if activation.is_full() || completed >= n_total {
|
|
return;
|
|
}
|
|
|
|
// Expect roughly 1 core of extra efficiency per worker activated in the
|
|
// last growth step (CPU-bound case); require at least CPU_SPAWN_THRESHOLD
|
|
// (20 %) of that expected gain before growing again. Scaling by the last
|
|
// step's size — not the cumulative total — keeps the bar meaningful
|
|
// regardless of how many workers are already active: growing by 8 should
|
|
// always take ~+1.6 cores to confirm, whether that's the 2nd growth step
|
|
// or the 20th.
|
|
let cpu_threshold = CPU_SPAWN_THRESHOLD * activation.last_step() as f64;
|
|
|
|
// Call both unconditionally (no `||` short-circuit): each sampler must
|
|
// advance its own window every tick, regardless of what the other one
|
|
// reports, or it would starve behind whichever signal fires first.
|
|
let cpu_wants_more = cpu_sample.do_i_activate(cpu_threshold);
|
|
let io_wants_more = io_sample.do_i_activate(IO_SPAWN_THRESHOLD * activation.last_step() as f64);
|
|
if !(cpu_wants_more || io_wants_more) {
|
|
return;
|
|
}
|
|
|
|
let grew = activation.grow(GROWTH_DIVISOR, n_total);
|
|
if grew > 0 {
|
|
debug!(
|
|
"activated {} worker(s) — {}/{} active",
|
|
grew,
|
|
activation.total(),
|
|
activation.max()
|
|
);
|
|
}
|
|
}
|