Skip to content
Closed
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
4 changes: 4 additions & 0 deletions differential-dataflow/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added

- `arrange_core_with_trace` and `reduce_with_tactic_and_trace`, which build the operator's trace with a caller-supplied factory in place of `Trace::new`, so a caller can keep a handle into the trace before `TraceAgent` takes ownership of it.

## [0.25.1](https://github.com/TimelyDataflow/differential-dataflow/compare/differential-dataflow-v0.25.0...differential-dataflow-v0.25.1) - 2026-07-15

### Other
Expand Down
26 changes: 24 additions & 2 deletions differential-dataflow/src/operators/arrange/arrangement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,8 @@
use timely::dataflow::operators::{Enter, vec::Map};
use timely::order::PartialOrder;
use timely::dataflow::{Scope, Stream};
use timely::dataflow::operators::generic::Operator;
use timely::dataflow::operators::generic::{Operator, OperatorInfo};
use timely::scheduling::Activator;
use timely::dataflow::channels::pact::{ParallelizationContract, Pipeline};
use timely::progress::Timestamp;
use timely::progress::Antichain;
Expand Down Expand Up @@ -186,7 +187,7 @@
while let Some(key) = cursor.get_key(batch) {
while let Some(val) = cursor.get_val(batch) {
for datum in logic(key, val) {
cursor.map_times(batch, |time, diff| {

Check warning on line 190 in differential-dataflow/src/operators/arrange/arrangement.rs

View workflow job for this annotation

GitHub Actions / Cargo clippy

`time` shadows a previous, unrelated binding
session.give((datum.clone(), <BatchCursor<Tr> as Cursor>::owned_time(time), <BatchCursor<Tr> as Cursor>::owned_diff(diff)));
});
}
Expand Down Expand Up @@ -310,6 +311,27 @@
name: &str,
batcher: impl FnOnce(Option<Logger>, usize) -> Ba,
) -> Arranged<'scope, TraceAgent<Tr>>
where
C: Container + Clone + 'static,
P: ParallelizationContract<Tr::Time, C>,
Ba: Batcher<C, Time = Tr::Time, Output: Into<Tr::Batch>> + 'static,
Tr: Trace+'static,
{
arrange_core_with_trace(stream, pact, name, batcher, Tr::new)
}

/// Arranges a stream of updates like [`arrange_core`], building the trace with `trace`.
///
/// `trace` receives the arguments [`Trace::new`] would, and is called once, on the operator's
/// construction. A caller that needs to reach the trace after [`TraceAgent`] takes ownership of it,
/// for example through a handle the trace shares, can obtain that handle here.
pub fn arrange_core_with_trace<'scope, P, C, Ba, Tr>(
stream: Stream<'scope, Tr::Time, C>,
pact: P,
name: &str,
batcher: impl FnOnce(Option<Logger>, usize) -> Ba,
trace: impl FnOnce(OperatorInfo, Option<Logger>, Option<Activator>) -> Tr,
) -> Arranged<'scope, TraceAgent<Tr>>
where
C: Container + Clone + 'static,
P: ParallelizationContract<Tr::Time, C>,
Expand Down Expand Up @@ -349,7 +371,7 @@
let mut capabilities = Antichain::<Capability<Tr::Time>>::new();

let activator = Some(scope.activator_for(std::rc::Rc::clone(&info.address)));
let mut empty_trace = Tr::new(info.clone(), logger.clone(), activator);
let mut empty_trace = trace(info.clone(), logger.clone(), activator);
// If there is default exertion logic set, install it.
if let Some(exert_logic) = scope.worker().config().get::<trace::ExertionLogic>("differential/default_exert_logic").cloned() {
empty_trace.set_exert_logic(exert_logic);
Expand Down
27 changes: 25 additions & 2 deletions differential-dataflow/src/operators/reduce.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,11 @@ use timely::progress::Timestamp;
use timely::dataflow::operators::Operator;
use timely::dataflow::operators::CapabilitySet;
use timely::dataflow::channels::pact::Pipeline;
use timely::dataflow::operators::generic::OperatorInfo;
use timely::scheduling::Activator;

use crate::operators::arrange::{Arranged, TraceAgent};
use crate::logging::Logger;
use crate::trace::{Span, ExertionLogic, Trace, TraceReader};

/// Sort and deduplicate a list. Shared by the cursor and proxy tactics, which
Expand Down Expand Up @@ -75,7 +78,27 @@ pub use crate::operators::cursor::reduce::reduce_trace;
/// `TraceReader` of its input and `Trace` of its output, never `Navigable`: it extracts batches via
/// `spans_through`, and building cursors over them (if that is how the reduce proceeds) is the
/// tactic's concern.
pub fn reduce_with_tactic<'scope, Tr1, Tr2, T>(trace: Arranged<'scope, Tr1>, name: &str, mut tactic: T) -> Arranged<'scope, TraceAgent<Tr2>>
pub fn reduce_with_tactic<'scope, Tr1, Tr2, T>(trace: Arranged<'scope, Tr1>, name: &str, tactic: T) -> Arranged<'scope, TraceAgent<Tr2>>
where
Tr1: TraceReader + 'static,
Tr2: Trace<Time = Tr1::Time> + 'static,
T: ReduceTactic<Tr1::Time, Tr1::Batch, Tr2::Batch> + 'static,
{
reduce_with_tactic_and_trace(trace, name, tactic, Tr2::new)
}

/// Drives a key-wise reduction like [`reduce_with_tactic`], building the output trace with
/// `output`.
///
/// `output` receives the arguments [`Trace::new`] would, and is called once, on the operator's
/// construction. See [`crate::operators::arrange::arrangement::arrange_core_with_trace`] for why a
/// caller would supply it.
pub fn reduce_with_tactic_and_trace<'scope, Tr1, Tr2, T>(
trace: Arranged<'scope, Tr1>,
name: &str,
mut tactic: T,
output: impl FnOnce(OperatorInfo, Option<Logger>, Option<Activator>) -> Tr2,
) -> Arranged<'scope, TraceAgent<Tr2>>
where
Tr1: TraceReader + 'static,
Tr2: Trace<Time = Tr1::Time> + 'static,
Expand All @@ -95,7 +118,7 @@ where
let logger = scope.worker().logger_for::<crate::logging::DifferentialEventBuilder>("differential/arrange").map(Into::into);

let activator = Some(scope.activator_for(std::rc::Rc::clone(&operator_info.address)));
let mut empty = Tr2::new(operator_info.clone(), logger.clone(), activator);
let mut empty = output(operator_info.clone(), logger.clone(), activator);
// If there is default exert logic set, install it.
if let Some(exert_logic) = scope.worker().config().get::<ExertionLogic>("differential/default_exert_logic").cloned() {
empty.set_exert_logic(exert_logic);
Expand Down
75 changes: 75 additions & 0 deletions differential-dataflow/tests/trace_factory.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
//! The `_with_trace` operator variants build their trace with the supplied factory.

use std::cell::Cell;
use std::rc::Rc;

use timely::dataflow::channels::pact::Exchange;
use timely::dataflow::operators::ToStream;

use differential_dataflow::AsCollection;
use differential_dataflow::hashable::Hashable;
use differential_dataflow::operators::arrange::arrangement::arrange_core_with_trace;
use differential_dataflow::operators::cursor::reduce::CursorTactic;
use differential_dataflow::operators::reduce::reduce_with_tactic_and_trace;
use differential_dataflow::trace::Trace;
use differential_dataflow::trace::implementations::{ValBatcher, ValBuilder, ValSpine};

type Spine = ValSpine<u64, u64, u64, isize>;

#[test]
fn arrange_builds_its_trace_with_the_factory() {
timely::example(|scope| {
let calls = Rc::new(Cell::new(0));
let built_for = Rc::new(Cell::new(None));
let (calls_in, built_for_in) = (Rc::clone(&calls), Rc::clone(&built_for));
let stream = vec![((1u64, 2u64), 0u64, 1isize)].to_stream(scope);
let exchange = Exchange::new(|update: &((u64, u64), u64, isize)| (update.0).0.hashed().into());
let arranged = arrange_core_with_trace::<_, _, _, Spine>(
stream,
exchange,
"Arrange",
ValBatcher::new,
move |info, logger, activator| {
calls_in.set(calls_in.get() + 1);
built_for_in.set(Some(info.global_id));
Spine::new(info, logger, activator)
},
);
assert_eq!(calls.get(), 1);
assert_eq!(built_for.get(), Some(arranged.trace.operator().global_id));
});
}

#[test]
fn reduce_builds_its_output_trace_with_the_factory() {
timely::example(|scope| {
let calls = Rc::new(Cell::new(0));
let built_for = Rc::new(Cell::new(None));
let (calls_in, built_for_in) = (Rc::clone(&calls), Rc::clone(&built_for));
let input = vec![((1u64, 2u64), 0u64, 1isize)]
.to_stream(scope)
.as_collection()
.arrange_by_key();
let tactic = CursorTactic::<_, _, ValBuilder<u64, u64, u64, isize>, _, _>::new(
|_key: &u64, input: &[(&u64, isize)], _output: &mut Vec<(u64, isize)>, change: &mut Vec<(u64, isize)>| {
change.extend(input.iter().map(|(v, r)| (**v, *r)));
},
|vec: &mut Vec<((u64, u64), u64, isize)>, key: &u64, upds: &mut Vec<(u64, u64, isize)>| {
vec.clear();
vec.extend(upds.drain(..).map(|(v, t, r)| ((*key, v), t, r)));
},
);
let reduced = reduce_with_tactic_and_trace::<_, Spine, _>(
input,
"Reduce",
tactic,
move |info, logger, activator| {
calls_in.set(calls_in.get() + 1);
built_for_in.set(Some(info.global_id));
Spine::new(info, logger, activator)
},
);
assert_eq!(calls.get(), 1);
assert_eq!(built_for.get(), Some(reduced.trace.operator().global_id));
});
}
Loading