Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions interactive/server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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, 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.
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
Expand Down
10 changes: 5 additions & 5 deletions interactive/server/src/loop_.rs
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,7 @@ pub fn run_worker(
worker: &mut Worker,
events: Option<Receiver<ControlEvent>>,
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
Expand Down Expand Up @@ -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<TailKey, Tail> = HashMap::new();
let mut responses: HashMap<u64, Sender<String>> = HashMap::new();
let mut next_token = 0u64;
Expand Down Expand Up @@ -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::<OuterTime, _, _>(|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() {
Expand Down
5 changes: 4 additions & 1 deletion interactive/server/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,9 @@ fn main() {
.unwrap_or_else(|_| "vec".to_string())
.parse::<RenderBackend>()
.unwrap_or_else(|error| panic!("DDIR_BACKEND: {error}"));
// 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::<ControlEvent>();
let (activation_tx, activation_rx) = sync_channel(1);
Expand All @@ -144,7 +147,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
Expand Down
113 changes: 111 additions & 2 deletions interactive/src/backend/corgi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<crate::backend::vec::Col<'s>>,
) -> Vec<crate::backend::vec::Col<'s>> {
) -> Vec<Collection<'s, Time, CC>> {
let corgi_imports: Vec<Collection<'s, Time, CC>> = crate::backend::vec::check_import_shapes(s, imports)
.into_iter()
.zip(&s.imports)
Expand Down Expand Up @@ -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<differential_dataflow::trace::chunk::ChunkSpine<CorgiChunk<u64, Diff>>>;

/// Arrange an export (already at host time) as a columnar trace.
pub fn arrange_export<'s>(c: Collection<'s, u64, CorgiContainer<u64, Diff>>) -> Arranged<'s, ExportTrace> {
arrange_core::<_, CorgiContainer<u64, Diff>, _, differential_dataflow::trace::chunk::ChunkSpine<CorgiChunk<u64, Diff>>>(
c.inner,
CorgiPact,
"CorgiExport",
ChunkBatcher::<CorgiChunker<u64, Diff>, _>::new,
)
}

/// The batches an export's arrangement has produced on this worker, not yet taken.
pub type ExportTap = std::rc::Rc<std::cell::RefCell<Vec<<ExportTrace as differential_dataflow::trace::TraceReader>::Batch>>>;

/// The updates in `batches`, as rows.
pub fn batch_rows(batches: Vec<<ExportTrace as differential_dataflow::trace::TraceReader>::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<<ExportTrace as differential_dataflow::trace::TraceReader>::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<differential_dataflow::trace::Span<u64, <ExportTrace as differential_dataflow::trace::TraceReader>::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<crate::backend::vec::Col<'s>>,
) -> Vec<crate::backend::vec::Col<'s>> {
render_tree_corgi(s, scope, depth, imports)
.into_iter()
.map(|c| {
c.inner
Expand Down
13 changes: 13 additions & 0 deletions interactive/src/corgi/container.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<R: Clone + 'static> differential_dataflow::collection::containers::Leave<crate::ir::Time, u64>
for CorgiContainer<crate::ir::Time, R>
{
type OuterContainer = CorgiContainer<u64, R>;
fn leave(self) -> CorgiContainer<u64, R> {
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 }
}
}
Loading
Loading