From 864d7a18bc7938f72aebc710594a916ce76ed18c Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Mon, 28 Sep 2026 20:59:32 -0400 Subject: [PATCH 1/2] DDIR server: columnar exports and export taps With the corgi backend, `Server::set_columnar_exports(true)` keeps a program's exports as corgi containers leaving its scope and arranges them as traces of corgi chunks, instead of converting each export to rows and arranging those. Readers convert to rows only as they read, through one path (`Server::published` / `Published::import_rows`): `snapshot` (peek), another program's `import`, `bind`, and the server binary's `tail`. Default off; the server binary turns it on with DDIR_COLUMNAR_EXPORTS=1. `Server::set_export_taps(true)` (with columnar exports) keeps each export's new batches per worker for an embedding program to drain as a change stream, `take_changes` or, in bounded row batches, `for_each_change_batch`, without a snapshot dataflow. Tests: columnar and row exports agree through snapshot, import and bind (including a bound counter); bounded draining matches a full drain and consumes changes once. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/server/README.md | 7 +- interactive/server/src/loop_.rs | 10 +- interactive/server/src/main.rs | 4 +- interactive/src/backend/corgi.rs | 113 +++++++++++- interactive/src/corgi/container.rs | 13 ++ interactive/src/server.rs | 172 +++++++++++++++++-- interactive/tests/server_change_batches.rs | 52 ++++++ interactive/tests/server_columnar_exports.rs | 72 ++++++++ 8 files changed, 414 insertions(+), 29 deletions(-) create mode 100644 interactive/tests/server_change_batches.rs create mode 100644 interactive/tests/server_columnar_exports.rs diff --git a/interactive/server/README.md b/interactive/server/README.md index aba99f27e..50aa43b14 100644 --- a/interactive/server/README.md +++ b/interactive/server/README.md @@ -12,14 +12,17 @@ defaults. Diagnostics are disabled by default so an idle server can park; `DDIR_DIAGNOSTICS=1` enables the diagnostics dataflow and its listener on `DDIR_DIAG_PORT` (default 51371). `DDIR_WORKERS` selects the number of worker threads (default 1), and `DDIR_BACKEND=vec|corgi` selects the renderer for installed -programs (default `vec`). +programs (default `vec`). With the Corgi backend, `DDIR_COLUMNAR_EXPORTS=1` keeps +exports as traces of Corgi chunks rather than rows (less memory); `peek`, `tail`, +`bind` and other programs' imports still read them as rows. One backend is selected for the whole server. The current registry is a transitional row-speaking bridge: inputs and imports convert from `Value` rows to Corgi columns at a program boundary, and exports convert back before they become shareable traces. A Corgi program stays columnar between those boundaries, but a production Corgi server should replace the bridge with native -columnar inputs and traces. +columnar inputs and traces. `DDIR_COLUMNAR_EXPORTS` is the first step on the export +side: the published trace is columnar, and rows are made only when it is read. Worker 0 admits one FIFO control stream and routes one ordered record to every worker. Small commands are replicated; framed input batches are partitioned by diff --git a/interactive/server/src/loop_.rs b/interactive/server/src/loop_.rs index 2d6af882c..0fb2ae8ec 100644 --- a/interactive/server/src/loop_.rs +++ b/interactive/server/src/loop_.rs @@ -102,6 +102,7 @@ pub fn run_worker( worker: &mut Worker, events: Option>, backend: RenderBackend, + columnar_exports: bool, ) { // Logging a park wakes the diagnostics dataflow, whose scheduling logs can // in turn wake it again. Keep idle servers genuinely idle unless an @@ -146,6 +147,7 @@ pub fn run_worker( let mut work_input = (worker.index() == 0).then_some(input); let mut server = Server::with_backend(backend); + server.set_columnar_exports(columnar_exports); let mut tails: HashMap = HashMap::new(); let mut responses: HashMap> = HashMap::new(); let mut next_token = 0u64; @@ -458,16 +460,14 @@ fn start_tail( return Err(format!("tail reqid {:?} is already active", reqid)); } let mut trace = server - .trace(name) + .published(name) .ok_or_else(|| format!("no trace {:?}", name))?; let dataflow_id = worker.next_dataflow_index(); let tag = reqid.to_string(); let mut probe = ProbeHandle::new(); let shutdown = worker.dataflow::(|scope| { - let (arranged, shutdown) = trace.import_core(scope.clone(), "TailImport"); - arranged - .as_collection(|key, val| (key.clone(), val.clone())) - .inner + let (rows, shutdown) = trace.import_rows(scope.clone(), "TailImport"); + rows .exchange(|_| 0u64) .inspect(move |((key, val), time, diff)| { if let Some(response) = response.as_ref() { diff --git a/interactive/server/src/main.rs b/interactive/server/src/main.rs index 7502d3a8d..76e1736d0 100644 --- a/interactive/server/src/main.rs +++ b/interactive/server/src/main.rs @@ -121,6 +121,8 @@ fn main() { .unwrap_or_else(|_| "vec".to_string()) .parse::() .unwrap_or_else(|error| panic!("DDIR_BACKEND: {error}")); + // DDIR_COLUMNAR_EXPORTS=1: corgi-backend exports stay columnar (less memory; readers see rows). + let columnar_exports = std::env::var("DDIR_COLUMNAR_EXPORTS").is_ok_and(|v| v != "0" && !v.is_empty()); let (event_tx, event_rx) = channel::(); let (activation_tx, activation_rx) = sync_channel(1); @@ -144,7 +146,7 @@ fn main() { } else { None }; - control_loop::run_worker(worker, events, backend); + control_loop::run_worker(worker, events, backend, columnar_exports); // The live server intentionally keeps installed dataflows and // diagnostics around. Once every worker observes shutdown, remove diff --git a/interactive/src/backend/corgi.rs b/interactive/src/backend/corgi.rs index 7c09c4e3a..29a5fa1e6 100644 --- a/interactive/src/backend/corgi.rs +++ b/interactive/src/backend/corgi.rs @@ -343,12 +343,12 @@ pub fn render_tree<'s>( /// converts back (`FromCorgi`). Signature-compatible with /// [`vec::render_tree`](crate::backend::vec::render_tree) (hence the `vec::Col` alias), so a /// row-speaking driver switches backends by switching this one call. -pub fn render_tree_rows<'s>( +pub fn render_tree_corgi<'s>( s: &st::Scope, scope: Scope<'s, Time>, depth: usize, imports: Vec>, -) -> Vec> { +) -> Vec> { let corgi_imports: Vec> = crate::backend::vec::check_import_shapes(s, imports) .into_iter() .zip(&s.imports) @@ -380,6 +380,115 @@ pub fn render_tree_rows<'s>( }) .collect(); render_tree(s, scope, depth, corgi_imports) +} + +/// A columnar export: the program's corgi collection, left to the host time and arranged as +/// corgi chunks. Readers convert to rows only when they read ([`export_rows`]). +pub type ExportTrace = TraceAgent>>; + +/// Arrange an export (already at host time) as a columnar trace. +pub fn arrange_export<'s>(c: Collection<'s, u64, CorgiContainer>) -> Arranged<'s, ExportTrace> { + arrange_core::<_, CorgiContainer, _, differential_dataflow::trace::chunk::ChunkSpine>>( + c.inner, + CorgiPact, + "CorgiExport", + ChunkBatcher::, _>::new, + ) +} + +/// The batches an export's arrangement has produced on this worker, not yet taken. +pub type ExportTap = std::rc::Rc::Batch>>>; + +/// The updates in `batches`, as rows. +pub fn batch_rows(batches: Vec<::Batch>) -> Vec<((Row, Row), u64, Diff)> { + let mut out = Vec::new(); + for batch in batches { + for ch in batch.chunks.iter().filter(|c| c.len() > 0) { + let c = CorgiContainer { + keys: recover_key(ch.keys()), + vals: ch.vals().clone(), + times: ch.times().clone(), + diffs: ch.diffs().to_vec(), + }; + out.extend(c.into_updates()); + } + } + out +} + +/// Decode and consume at most `rows_per_batch` export rows at a time. +/// Row payloads (including nested lists) can still vary in size. +pub fn for_each_batch_rows( + batches: Vec<::Batch>, + rows_per_batch: usize, + mut consume: impl FnMut(Vec<((Row, Row), u64, Diff)>), +) { + assert!(rows_per_batch > 0, "export row batch size must be positive"); + let mut indices = Vec::new(); + for batch in batches { + for ch in batch.chunks.iter().filter(|c| c.len() > 0) { + let keys = recover_key(ch.keys()); + for start in (0..ch.len()).step_by(rows_per_batch) { + let end = start.saturating_add(rows_per_batch).min(ch.len()); + indices.clear(); + indices.extend(start..end); + let c = CorgiContainer { + keys: corgi::arrange::gather(&keys, &indices), + vals: corgi::arrange::gather(ch.vals(), &indices), + times: ch.times().gather(&indices), + diffs: ch.diffs()[start..end].to_vec(), + }; + consume(c.into_updates()); + } + } + } +} + +/// Record each batch the arrangement emits into `tap`, passing the stream through. +pub fn tap_export<'s>(a: &Arranged<'s, ExportTrace>, tap: ExportTap) -> timely::dataflow::Stream<'s, u64, Vec::Batch>>> { + a.stream.clone().unary(Pipeline, "ExportTap", move |_, _| { + move |input, output| { + input.for_each(|cap, data| { + tap.borrow_mut().extend(data.iter().filter_map(|b| b.inner.clone())); + output.session(&cap).give_container(data); + }); + } + }) +} + +/// The rows of an imported columnar export, as `((key, val), time, diff)` updates. +pub fn export_rows<'s>( + a: Arranged<'s, ExportTrace>, +) -> timely::dataflow::Stream<'s, u64, Vec<((Row, Row), u64, Diff)>> { + a.stream.unary(Pipeline, "ExportRows", |_, _| { + |input, output| { + input.for_each(|cap, data| { + let mut session = output.session(&cap); + for batch in data.iter() { + let Some(payload) = batch.inner.as_ref() else { continue }; + for ch in payload.chunks.iter().filter(|c| c.len() > 0) { + let c = CorgiContainer { + keys: recover_key(ch.keys()), + vals: ch.vals().clone(), + times: ch.times().clone(), + diffs: ch.diffs().to_vec(), + }; + session.give_container(&mut c.into_updates()); + } + } + }); + } + }) +} + +/// [`render_tree_corgi`] with each export converted back to rows (`FromCorgi`). +pub fn render_tree_rows<'s>( + s: &st::Scope, + scope: Scope<'s, Time>, + depth: usize, + imports: Vec>, +) -> Vec> { + render_tree_corgi(s, scope, depth, imports) .into_iter() .map(|c| { c.inner diff --git a/interactive/src/corgi/container.rs b/interactive/src/corgi/container.rs index 2ae94094b..af0a2c410 100644 --- a/interactive/src/corgi/container.rs +++ b/interactive/src/corgi/container.rs @@ -147,3 +147,16 @@ where CorgiContainer { keys, vals, times: self.times, diffs } } } + + +/// Leaving the program's iterative scope: each time keeps its outer coordinate (lane 0). +impl differential_dataflow::collection::containers::Leave + for CorgiContainer +{ + type OuterContainer = CorgiContainer; + fn leave(self) -> CorgiContainer { + let len = self.times.len(); + let lane = if self.times.width() == 0 { vec![0; len] } else { self.times.lane(0) }; + CorgiContainer { keys: self.keys, vals: self.vals, times: ColTimes::from_raw_lanes(vec![lane], len), diffs: self.diffs } + } +} diff --git a/interactive/src/server.rs b/interactive/src/server.rs index 62c242b53..fb285943e 100644 --- a/interactive/src/server.rs +++ b/interactive/src/server.rs @@ -106,6 +106,38 @@ impl std::str::FromStr for RenderBackend { /// arranged by key at the host time so any later install can `import` it. pub type ServerTrace = TraceAgent>; +/// A published export as its readers see it: a row trace, or with columnar exports +/// ([`Server::set_columnar_exports`]) a trace of corgi chunks. Either reads as rows. +#[derive(Clone)] +pub enum Published { + Rows(ServerTrace), + Columnar(crate::backend::corgi::ExportTrace), +} + +impl Published { + /// Import into `scope` as a stream of `((key, val), time, diff)` rows, with the + /// import's shutdown button. A columnar trace converts to rows as it is read. + pub fn import_rows<'s>( + &mut self, + scope: timely::dataflow::Scope<'s, OuterTime>, + name: &str, + ) -> ( + timely::dataflow::Stream<'s, OuterTime, Vec<((Value, Value), OuterTime, Diff)>>, + ShutdownButton>, + ) { + match self { + Published::Rows(trace) => { + let (arranged, shutdown) = trace.import_core(scope, name); + (arranged.as_collection(|k, v| (k.clone(), v.clone())).inner, shutdown) + } + Published::Columnar(trace) => { + let (arranged, shutdown) = trace.import_core(scope, name); + (crate::backend::corgi::export_rows(arranged), shutdown) + } + } + } +} + /// An input handle into an installed program's positional `input N`. type ServerInput = InputSession; @@ -353,6 +385,16 @@ struct Binding { pub struct Server { /// Published export name -> shareable trace. traces: HashMap, + /// With columnar exports: published export name -> trace of corgi chunks. Readers + /// (imports, binds, snapshots, subscriptions) see rows, converted as they read. + ctraces: HashMap, + /// With export taps: each columnar export's new batches on this worker, drained by + /// `take_changes` / `for_each_change_batch`: a standing change stream without a snapshot. + taps: HashMap, + /// `set_columnar_exports`. + columnar_exports: bool, + /// `set_export_taps`. + export_taps: bool, /// Installed program name -> its handles and lifecycle bookkeeping. programs: HashMap, /// Trace name -> number of installed programs importing it (the drop gate). @@ -376,6 +418,10 @@ impl Server { pub fn with_backend(backend: RenderBackend) -> Self { Server { traces: HashMap::new(), + ctraces: HashMap::new(), + taps: HashMap::new(), + columnar_exports: false, + export_taps: false, programs: HashMap::new(), importers: HashMap::new(), bindings: Vec::new(), @@ -389,11 +435,37 @@ impl Server { self.epoch } + /// Whether `name` (canonical) is published, as a row or a columnar trace. + fn is_published(&self, name: &str) -> bool { + self.traces.contains_key(name) || self.ctraces.contains_key(name) + } + + /// Keep corgi-backend exports columnar (default off): a program's exports leave its scope as + /// corgi containers and are arranged as corgi chunks, and readers convert to rows only when + /// they read (`snapshot`, imports, binds). Affects programs installed afterwards. + pub fn set_columnar_exports(&mut self, on: bool) { + self.columnar_exports = on; + } + + /// With columnar exports, also keep each export's new batches for [`Server::take_changes`] + /// and [`Server::for_each_change_batch`] (default off). Affects programs installed afterwards. + pub fn set_export_taps(&mut self, on: bool) { + self.export_taps = on; + } + /// Clone a trace reader for a transient peek or subscription dataflow. pub fn trace(&self, name: &str) -> Option { self.traces.get(&canonical_source_name(name)).cloned() } + /// A reader for published trace `name`, row or columnar, for a transient peek or + /// subscription dataflow. + pub fn published(&self, name: &str) -> Option { + let name = canonical_source_name(name); + self.traces.get(&name).cloned().map(Published::Rows) + .or_else(|| self.ctraces.get(&name).cloned().map(Published::Columnar)) + } + /// Return registry state without coupling a caller to stdout formatting. pub fn program_info(&self) -> Vec { let mut result: Vec<_> = self @@ -455,7 +527,7 @@ impl Server { for imp in &prog.root.imports { if let st::Source::Trace(t) = &imp.from { let key = canonical_source_name(t); - if !self.traces.contains_key(&key) { + if !self.is_published(&key) { if key != "clock" && Recipe::parse(&key).is_none() { return Err(format!( "program {:?} imports unknown trace {:?}; install its producer first", @@ -467,7 +539,7 @@ impl Server { } } for e in &prog.root.exports { - if self.traces.contains_key(&e.name) || generated.contains(&e.name) { + if self.is_published(&e.name) || generated.contains(&e.name) { return Err(format!("export name {:?} is already published; choose another name or drop its producer", e.name)); } } @@ -498,11 +570,15 @@ impl Server { let root = &prog.root; let traces = &mut self.traces; let backend = self.backend; + let columnar = backend == RenderBackend::Corgi && self.columnar_exports; + let tapped = columnar && self.export_taps; + let ctraces = &mut self.ctraces; // The id this dataflow will get; captured so `drop` can remove it. let dataflow_id = worker.next_dataflow_index(); - let (published, inputs): (Vec<(String, ServerTrace)>, Vec<(usize, ServerInput)>) = + #[allow(clippy::type_complexity)] + let (published, cpublished, inputs): (Vec<(String, ServerTrace)>, Vec<(String, crate::backend::corgi::ExportTrace, Option)>, Vec<(usize, ServerInput)>) = worker.dataflow::(|outer| { let mut inputs: Vec<(usize, ServerInput)> = Vec::new(); @@ -519,6 +595,11 @@ impl Server { st::Source::Trace(t) => { // The first binding point: resolve a named trace by importing it. let key = canonical_source_name(t); + // A columnar export is read as rows by its importers. + if let Some(ctrace) = ctraces.get_mut(&key) { + let rows = crate::backend::corgi::export_rows(ctrace.import(outer.clone())); + return differential_dataflow::AsCollection::as_collection(rows); + } let arranged = traces .get_mut(&key) .expect("validated above") @@ -531,21 +612,41 @@ impl Server { // Render the program body in its own iterative scope, then bring // every export back out to the host time. - let leaved: Vec> = outer + let (leaved, cleaved) = outer .iterative::, _, _>(|inner| { let entered: Vec<_> = outer_cols.iter().map(|c| c.clone().enter(inner)).collect(); + if columnar { + let exports = crate::backend::corgi::render_tree_corgi(root, inner.clone(), 0, entered); + return (Vec::new(), exports.into_iter().map(|c| c.leave(outer)).collect::>()); + } let exports = match backend { RenderBackend::Vec => render_tree(root, inner.clone(), 0, entered), RenderBackend::Corgi => { render_tree_rows(root, inner.clone(), 0, entered) } }; - exports + (exports .into_iter() .map(|c| c.leave(outer)) - .collect::>() + .collect::>(), Vec::new()) }); + let cpublished: Vec<_> = root + .exports + .iter() + .zip(cleaved) + .map(|(e, col)| { + use timely::dataflow::operators::Probe; + let arranged = crate::backend::corgi::arrange_export(col); + let tap: Option = tapped.then(Default::default); + let stream = match &tap { + Some(t) => crate::backend::corgi::tap_export(&arranged, t.clone()), + None => arranged.stream.clone(), + }; + stream.probe_with(&probe); + (e.name.clone(), arranged.trace, tap) + }) + .collect(); // The second binding point: probe and publish each export's trace. let published: Vec<(String, ServerTrace)> = root @@ -560,12 +661,16 @@ impl Server { }) .collect(); - (published, inputs) + (published, cpublished, inputs) }); for (export_name, trace) in published { self.traces.insert(export_name, trace); } + for (export_name, trace, tap) in cpublished { + if let Some(tap) = tap { self.taps.insert(export_name.clone(), tap); } + self.ctraces.insert(export_name, trace); + } for t in &import_names { *self.importers.entry(t.clone()).or_insert(0) += 1; } @@ -862,7 +967,7 @@ impl Server { ) -> Result<(), String> { let source = canonical_source_name(trace); let target = prog.to_string(); - if !self.traces.contains_key(&source) { + if !self.is_published(&source) { return Err(format!("no trace {:?}", source)); } let installed = self @@ -890,11 +995,11 @@ impl Server { let buffer_in = buffer.clone(); let mut probe = ProbeHandle::new(); let dataflow_id = worker.next_dataflow_index(); - let trace_handle = self.traces.get_mut(&source).expect("checked above"); + let mut trace = self.published(&source).expect("checked above"); let shutdown = worker.dataflow::(|scope| { - let (arranged, shutdown) = trace_handle.import_core(scope.clone(), "BindImport"); - arranged - .as_collection(|k, v| (k.clone(), v.clone())) + use timely::dataflow::operators::{Inspect, Probe}; + let (rows, shutdown) = trace.import_rows(scope.clone(), "BindImport"); + rows .inspect(move |((key, val), _time, diff)| { buffer_in .borrow_mut() @@ -970,18 +1075,15 @@ impl Server { let name = canonical_source_name(name); let epoch = self.epoch; - let mut trace = self - .trace(&name) - .ok_or_else(|| format!("no trace {:?}", name))?; + let mut trace = self.published(&name).ok_or_else(|| format!("no trace {:?}", name))?; let acc: Rc>> = Rc::new(RefCell::new(HashMap::new())); let acc_in = acc.clone(); let mut probe = ProbeHandle::new(); let id = worker.next_dataflow_index(); worker.dataflow::(|scope| { - trace - .import(scope.clone()) - .as_collection(|k, v| (k.clone(), v.clone())) - .inner + // The dataflow is dropped once drained, which releases the import. + let (rows, _shutdown) = trace.import_rows(scope.clone(), "SnapshotImport"); + rows .exchange(|_| 0u64) .inspect(move |((k, v), t, d)| { if *t < epoch { @@ -1008,6 +1110,32 @@ impl Server { Ok(rows) } + /// The changes to export `name` that this worker's share of its arrangement has produced + /// since the last call, as rows (with [`Server::set_export_taps`]). Each worker drains its own. + pub fn take_changes(&mut self, name: &str) -> Vec<((Value, Value), OuterTime, Diff)> { + match self.taps.get(name) { + Some(tap) => crate::backend::corgi::batch_rows(std::mem::take(&mut *tap.borrow_mut())), + None => Vec::new(), + } + } + + /// Drain this worker's tapped export changes in row batches no larger than + /// `rows_per_batch`. Callbacks run in export order; rows are consumed once. + /// Unlike `take_changes`, this does not materialize the entire export as rows. + /// The limit counts rows, not bytes (nested payloads may be large). + pub fn for_each_change_batch( + &mut self, + name: &str, + rows_per_batch: usize, + consume: impl FnMut(Vec<((Value, Value), OuterTime, Diff)>), + ) { + assert!(rows_per_batch > 0, "export row batch size must be positive"); + if let Some(tap) = self.taps.get(name) { + let batches = std::mem::take(&mut *tap.borrow_mut()); + crate::backend::corgi::for_each_batch_rows(batches, rows_per_batch, consume); + } + } + /// Drop installed program `name`, releasing its dataflow immediately. /// /// Refuses (changing nothing) if any trace the program publishes still has a @@ -1046,6 +1174,8 @@ impl Server { } for ex in &installed.exports { self.traces.remove(ex); + self.ctraces.remove(ex); + self.taps.remove(ex); } let id = installed.dataflow_id; // Drop the input handles first (closes the inputs while the operators @@ -1167,6 +1297,10 @@ impl Server { trace.set_logical_compaction(frontier.borrow()); trace.set_physical_compaction(frontier.borrow()); } + for trace in self.ctraces.values_mut() { + trace.set_logical_compaction(frontier.borrow()); + trace.set_physical_compaction(frontier.borrow()); + } } } diff --git a/interactive/tests/server_change_batches.rs b/interactive/tests/server_change_batches.rs new file mode 100644 index 000000000..1ee83ed75 --- /dev/null +++ b/interactive/tests/server_change_batches.rs @@ -0,0 +1,52 @@ +//! Bounded export consumption must preserve inserts/retractions and nested rows. +use interactive::ir::Value; +use interactive::server::{InputUpdate, RenderBackend, Server}; +use interactive::{lower, parse}; + +#[test] +fn bounded_exports_match_full_drain_and_are_consumed_once() { + timely::execute_directly(|worker| { + let mut program = lower::lower_tree(parse::pipe::parse(r#" + type Three = A u64 | B u64 | C u64; + let pairs = input 0 | key($0[0] ; $0[1]); + export "plain" = pairs; + export "lists" = pairs | collect; + export "tagged" = pairs | map($0 ; variant(Three, 2, $1[0])); + "#)); + program.optimize(); + let mut full = Server::with_backend(RenderBackend::Corgi); + let mut bounded = Server::with_backend(RenderBackend::Corgi); + for server in [&mut full, &mut bounded] { + server.set_columnar_exports(true); + server.set_export_taps(true); + } + full.install(worker, "full", &program).unwrap(); + bounded.install(worker, "bounded", &program).unwrap(); + for (epoch, limit) in [1, 7, 128, 4096].into_iter().enumerate() { + let diff = if epoch % 2 == 0 { 1 } else { -1 }; + let input: Vec<_> = (0..257).map(|i| InputUpdate { + key: Value::Tuple(vec![Value::Int(i % 11), Value::Int(i - 128)]), + val: Value::unit(), diff, + }).collect(); + full.feed_batch("full", 0, input.clone()).unwrap(); + bounded.feed_batch("bounded", 0, input).unwrap(); + full.tick(worker); + bounded.tick(worker); + for name in ["plain", "lists", "tagged"] { + let mut expected = full.take_changes(name); + let mut actual = Vec::new(); + bounded.for_each_change_batch(name, limit, |rows| { + assert!(!rows.is_empty() && rows.len() <= limit); + actual.extend(rows); + }); + assert!(!expected.is_empty(), "{name}, epoch {epoch}"); + expected.sort(); + actual.sort(); + assert_eq!(actual, expected, "{name}, epoch {epoch}, limit {limit}"); + bounded.for_each_change_batch(name, limit, |_| panic!("changes drained twice")); + assert!(bounded.take_changes(name).is_empty()); + } + } + bounded.for_each_change_batch("missing", 7, |_| panic!("missing export has changes")); + }); +} diff --git a/interactive/tests/server_columnar_exports.rs b/interactive/tests/server_columnar_exports.rs new file mode 100644 index 000000000..8d0f4d13f --- /dev/null +++ b/interactive/tests/server_columnar_exports.rs @@ -0,0 +1,72 @@ +//! Columnar exports (`Server::set_columnar_exports`) read the same as row exports: +//! through `snapshot`, through another program's `import`, and through `bind`. + +use interactive::ir::Value; +use interactive::server::{RenderBackend, Server}; +use interactive::{lower, parse}; + +fn tup(fields: &[i64]) -> Value { + Value::Tuple(fields.iter().map(|&n| Value::Int(n)).collect()) +} + +fn install(server: &mut Server, worker: &mut timely::worker::Worker, name: &str, src: &str) { + let mut program = lower::lower_tree(parse::pipe::parse(src)); + program.optimize(); + server.install(worker, name, &program).unwrap(); +} + +const PRODUCER: &str = r#" + type Three = A u64 | B u64 | C u64; + let pairs = input 0 | key($0[0] ; $0[1]); + export "plain" = pairs; + export "lists" = pairs | collect; + export "tagged" = pairs | map($0 ; variant(Three, 2, $1[0])); +"#; + +const CONSUMER: &str = r#" + let p = import "plain"; + export "shifted" = p | map($0 ; $1[0] + 1000); +"#; + +const COUNTER: &str = r#" + let seed = input 0; + let feedback = input 1; + let state = seed + feedback; + export "count" = state; + export "next" = (state | map($0[0] + 1 ;)) + (seed | negate); +"#; + +/// Everything the server can read back, after the same inputs and ticks. +fn run(worker: &mut timely::worker::Worker, columnar: bool) -> Vec> { + let mut server = Server::with_backend(RenderBackend::Corgi); + server.set_columnar_exports(columnar); + install(&mut server, worker, "producer", PRODUCER); + install(&mut server, worker, "consumer", CONSUMER); + install(&mut server, worker, "counter", COUNTER); + server.feed("counter", 0, tup(&[0]), Value::unit(), None, 1).unwrap(); + server.bind(worker, "next", "counter", 1).unwrap(); + let mut seen = Vec::new(); + for epoch in 0..4i64 { + for i in 0..40i64 { + let diff = if (i + epoch) % 3 == 0 { -1 } else { 1 }; + server.feed("producer", 0, tup(&[i % 7, i * epoch]), Value::unit(), None, diff).unwrap(); + } + server.tick(worker); + for name in ["plain", "lists", "tagged", "shifted", "count"] { + seen.push(server.snapshot(worker, name).unwrap()); + } + } + seen +} + +#[test] +fn columnar_exports_read_like_row_exports() { + timely::execute_directly(|worker| { + let rows = run(worker, false); + let columnar = run(worker, true); + assert!(rows.iter().any(|s| !s.is_empty())); + // The counter advanced through the binding in both. + assert_eq!(rows.last(), Some(&vec![(tup(&[3]), Value::unit(), 1)])); + assert_eq!(rows, columnar); + }); +} From 4f02e0cbf8ba5871924811dc3e6029d52b878e09 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Mon, 28 Sep 2026 21:33:27 -0400 Subject: [PATCH 2/2] DDIR server: columnar exports on by default for the corgi backend Nothing measured is worse with them on, and the tour's peak falls 32-50%. set_columnar_exports(false), or DDIR_COLUMNAR_EXPORTS=0 for the server binary, keeps row traces. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/server/README.md | 10 +++++----- interactive/server/src/main.rs | 5 +++-- interactive/src/server.rs | 4 ++-- 3 files changed, 10 insertions(+), 9 deletions(-) diff --git a/interactive/server/README.md b/interactive/server/README.md index 50aa43b14..240e11531 100644 --- a/interactive/server/README.md +++ b/interactive/server/README.md @@ -12,17 +12,17 @@ defaults. Diagnostics are disabled by default so an idle server can park; `DDIR_DIAGNOSTICS=1` enables the diagnostics dataflow and its listener on `DDIR_DIAG_PORT` (default 51371). `DDIR_WORKERS` selects the number of worker threads (default 1), and `DDIR_BACKEND=vec|corgi` selects the renderer for installed -programs (default `vec`). With the Corgi backend, `DDIR_COLUMNAR_EXPORTS=1` keeps -exports as traces of Corgi chunks rather than rows (less memory); `peek`, `tail`, -`bind` and other programs' imports still read them as rows. +programs (default `vec`). With the Corgi backend, exports are kept as traces of +Corgi chunks rather than rows (less memory); `peek`, `tail`, `bind` and other +programs' imports still read them as rows. `DDIR_COLUMNAR_EXPORTS=0` keeps row traces. One backend is selected for the whole server. The current registry is a transitional row-speaking bridge: inputs and imports convert from `Value` rows to Corgi columns at a program boundary, and exports convert back before they become shareable traces. A Corgi program stays columnar between those boundaries, but a production Corgi server should replace the bridge with native -columnar inputs and traces. `DDIR_COLUMNAR_EXPORTS` is the first step on the export -side: the published trace is columnar, and rows are made only when it is read. +columnar inputs and traces. Columnar exports are the first step on the export side: +the published trace is columnar, and rows are made only when it is read. Worker 0 admits one FIFO control stream and routes one ordered record to every worker. Small commands are replicated; framed input batches are partitioned by diff --git a/interactive/server/src/main.rs b/interactive/server/src/main.rs index 76e1736d0..3498df951 100644 --- a/interactive/server/src/main.rs +++ b/interactive/server/src/main.rs @@ -121,8 +121,9 @@ fn main() { .unwrap_or_else(|_| "vec".to_string()) .parse::() .unwrap_or_else(|error| panic!("DDIR_BACKEND: {error}")); - // DDIR_COLUMNAR_EXPORTS=1: corgi-backend exports stay columnar (less memory; readers see rows). - let columnar_exports = std::env::var("DDIR_COLUMNAR_EXPORTS").is_ok_and(|v| v != "0" && !v.is_empty()); + // Corgi-backend exports stay columnar (less memory; readers see rows); DDIR_COLUMNAR_EXPORTS=0 + // keeps the row traces instead. + let columnar_exports = std::env::var("DDIR_COLUMNAR_EXPORTS").map_or(true, |v| v != "0"); let (event_tx, event_rx) = channel::(); let (activation_tx, activation_rx) = sync_channel(1); diff --git a/interactive/src/server.rs b/interactive/src/server.rs index fb285943e..499bdc742 100644 --- a/interactive/src/server.rs +++ b/interactive/src/server.rs @@ -420,7 +420,7 @@ impl Server { traces: HashMap::new(), ctraces: HashMap::new(), taps: HashMap::new(), - columnar_exports: false, + columnar_exports: true, export_taps: false, programs: HashMap::new(), importers: HashMap::new(), @@ -440,7 +440,7 @@ impl Server { self.traces.contains_key(name) || self.ctraces.contains_key(name) } - /// Keep corgi-backend exports columnar (default off): a program's exports leave its scope as + /// Keep corgi-backend exports columnar (default on): a program's exports leave its scope as /// corgi containers and are arranged as corgi chunks, and readers convert to rows only when /// they read (`snapshot`, imports, binds). Affects programs installed afterwards. pub fn set_columnar_exports(&mut self, on: bool) {