Skip to content
Draft
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

- `Agent` trait over the shared reader an arranging operator returns, implemented by `TraceAgent`, and `arrange_core_with_agent`, `reduce_with_tactic_and_agent`, and `arrange_from_upsert_with_agent`, which return a caller-chosen agent. `Agent::new` receives the trace by value, so an agent can keep state the trace shares 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
22 changes: 22 additions & 0 deletions differential-dataflow/src/operators/arrange/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,28 @@
fn map_spans<F: FnMut(&Span<Tr::Time, Tr::Batch>)>(&self, f: F) { self.trace.borrow().trace.map_spans(f) }
}

/// A shared reader of a trace, constructed by the operator that maintains the trace.
///
/// The arranging operators (`arrange_core_with_agent`, `reduce_with_tactic_and_agent`,
/// `arrange_from_upsert_with_agent`) build their trace, hand it to `Agent::new`, keep the writer,
/// and return the agent in the resulting `Arranged`. The operator also reads through its own copy of
/// the agent, so an implementation must honour the `TraceReader` contract of the trace it shares.
pub trait Agent: TraceReader + Sized {
/// The trace the agent shares.
type Trace: Trace<Time = Self::Time, Batch = Self::Batch>;
/// Takes ownership of `trace` and returns a reader and the writer that feeds it.
///
/// Called once, on the operator's construction.
fn new(trace: Self::Trace, operator: OperatorInfo, logging: Option<crate::logging::Logger>) -> (Self, TraceWriter<Self::Trace>);
}

impl<Tr: Trace> Agent for TraceAgent<Tr> {
type Trace = Tr;
fn new(trace: Tr, operator: OperatorInfo, logging: Option<crate::logging::Logger>) -> (Self, TraceWriter<Tr>) {
TraceAgent::new(trace, operator, logging)
}
}

impl<Tr: TraceReader> TraceAgent<Tr> {
/// Creates a new agent from a trace reader.
pub fn new(trace: Tr, operator: OperatorInfo, logging: Option<crate::logging::Logger>) -> (Self, TraceWriter<Tr>)
Expand Down Expand Up @@ -283,7 +305,7 @@
let activator = scope.activator_for(Rc::clone(&info.address));
let queue = self.new_listener(activator);

let activator = scope.activator_for(info.address);

Check warning on line 308 in differential-dataflow/src/operators/arrange/agent.rs

View workflow job for this annotation

GitHub Actions / Cargo clippy

`activator` shadows a previous, unrelated binding
*shutdown_button_ref = Some(ShutdownButton::new(Rc::clone(&capabilities), activator));

capabilities.borrow_mut().as_mut().unwrap().insert(capability);
Expand Down Expand Up @@ -414,7 +436,7 @@
let activator = scope.activator_for(Rc::clone(&info.address));
let queue = self.new_listener(activator);

let activator = scope.activator_for(info.address);

Check warning on line 439 in differential-dataflow/src/operators/arrange/agent.rs

View workflow job for this annotation

GitHub Actions / Cargo clippy

`activator` shadows a previous, unrelated binding
*shutdown_button_ref = Some(ShutdownButton::new(Rc::clone(&capabilities), activator));

capabilities.borrow_mut().as_mut().unwrap().insert(capability);
Expand Down
31 changes: 24 additions & 7 deletions differential-dataflow/src/operators/arrange/arrangement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@

use trace::wrappers::enter::{TraceEnter, enter_span};

use super::TraceAgent;
use super::{Agent, TraceAgent};

/// An arranged collection of `(K,V)` values.
///
Expand Down Expand Up @@ -186,7 +186,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 189 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 @@ -315,6 +315,23 @@
P: ParallelizationContract<Tr::Time, C>,
Ba: Batcher<C, Time = Tr::Time, Output: Into<Tr::Batch>> + 'static,
Tr: Trace+'static,
{
arrange_core_with_agent::<P, C, Ba, TraceAgent<Tr>>(stream, pact, name, batcher)
}

/// Arranges a stream of updates like [`arrange_core`], sharing the trace through the agent `A`.
pub fn arrange_core_with_agent<'scope, P, C, Ba, A>(
stream: Stream<'scope, A::Time, C>,
pact: P,
name: &str,
batcher: impl FnOnce(Option<Logger>, usize) -> Ba,
) -> Arranged<'scope, A>
where
C: Container + Clone + 'static,
P: ParallelizationContract<A::Time, C>,
Ba: Batcher<C, Time = A::Time, Output: Into<A::Batch>> + 'static,
A: Agent + 'static,
A::Trace: 'static,
{
// The `Arrange` operator is tasked with reacting to an advancing input
// frontier by producing the sequence of batches whose lower and upper
Expand All @@ -331,7 +348,7 @@
// held by the batcher, which may prevents the operator from sending an
// empty batch.

let mut reader: Option<TraceAgent<Tr>> = None;
let mut reader: Option<A> = None;

// fabricate a data-parallel operator that holds capabilities and consults its input frontier.
let reader_ref = &mut reader;
Expand All @@ -346,21 +363,21 @@
let mut batcher = batcher(logger.clone(), info.global_id);

// Capabilities for the lower envelope of updates in `batcher`.
let mut capabilities = Antichain::<Capability<Tr::Time>>::new();
let mut capabilities = Antichain::<Capability<A::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 = A::Trace::new(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);
}

let (reader_local, mut writer) = TraceAgent::new(empty_trace, info, logger);
let (reader_local, mut writer) = A::new(empty_trace, info, logger);

*reader_ref = Some(reader_local);

// Initialize to the minimal input frontier.
let mut prev_frontier = Antichain::from_elem(Tr::Time::minimum());
let mut prev_frontier = Antichain::from_elem(<A::Time as Timestamp>::minimum());

move |(input, frontier), output| {

Expand Down Expand Up @@ -418,7 +435,7 @@
let description = Description::new(
prev_frontier.clone(),
frontier.frontier().to_owned(),
Antichain::from_elem(Tr::Time::minimum()),
Antichain::from_elem(<A::Time as Timestamp>::minimum()),
);
let (extracted, retained) = batcher.extract(frontier.frontier());

Expand Down
2 changes: 1 addition & 1 deletion differential-dataflow/src/operators/arrange/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,6 @@ pub mod arrangement;
pub mod upsert;

pub use self::writer::TraceWriter;
pub use self::agent::{TraceAgent, ShutdownButton};
pub use self::agent::{Agent, TraceAgent, ShutdownButton};

pub use self::arrangement::Arranged;
45 changes: 33 additions & 12 deletions differential-dataflow/src/operators/arrange/upsert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,12 +108,12 @@ use timely::progress::{Antichain, Timestamp};
use timely::dataflow::operators::Capability;

use crate::operators::arrange::arrangement::Arranged;
use crate::trace::{self, BatchCursor, BatchDiff, Builder, Cursor, Description, Navigable, Trace, TraceReader};
use crate::trace::{self, BatchCursor, BatchDiff, Builder, Cursor, Description, Navigable, Trace};
use crate::{ExchangeData, Hashable};

use crate::trace::implementations::containers::BatchContainer;

use super::TraceAgent;
use super::{Agent, TraceAgent};

/// Arrange data from a stream of keyed upserts.
///
Expand All @@ -140,14 +140,35 @@ where
>,
Bu: Builder<Time=Tr::Time, Input = Vec<((K, V), Tr::Time, BatchDiff<Tr>)>, Output: Into<Tr::Batch>>,
{
let mut reader: Option<TraceAgent<Tr>> = None;
arrange_from_upsert_with_agent::<Bu, TraceAgent<Tr>, K, V>(stream, name)
}

/// Arranges data from a stream of keyed upserts like [`arrange_from_upsert`], sharing the trace
/// through the agent `A`.
pub fn arrange_from_upsert_with_agent<'scope, Bu, A, K, V>(
stream: Stream<'scope, A::Time, Vec<(K, Option<V>, A::Time)>>,
name: &str,
) -> Arranged<'scope, A>
where
K: ExchangeData+Hashable+std::hash::Hash,
V: ExchangeData,
A: Agent<Batch: Navigable, Time: TotalOrder+ExchangeData> + Clone + 'static,
A::Trace: 'static,
for<'a> BatchCursor<A>: Cursor<
Key<'a> = &'a K,
Val<'a> = &'a V,
Diff=isize,
>,
Bu: Builder<Time=A::Time, Input = Vec<((K, V), A::Time, BatchDiff<A>)>, Output: Into<A::Batch>>,
{
let mut reader: Option<A> = None;

// fabricate a data-parallel operator that holds capabilities and consults its input frontier.
let stream = {

let reader = &mut reader;

let exchange = Exchange::new(move |update: &(K,Option<V>,Tr::Time)| (update.0).hashed().into());
let exchange = Exchange::new(move |update: &(K,Option<V>,A::Time)| (update.0).hashed().into());

let scope = stream.scope();
stream.unary_frontier(exchange, name, move |_capability, info| {
Expand All @@ -156,24 +177,24 @@ where
let logger = scope.worker().logger_for::<crate::logging::DifferentialEventBuilder>("differential/arrange").map(Into::into);

// Tracks the lower envelope of times in `priority_queue`.
let mut capabilities = Antichain::<Capability<Tr::Time>>::new();
let mut capabilities = Antichain::<Capability<A::Time>>::new();
// Form the trace we will both use internally and publish.
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 = A::Trace::new(info.clone(), logger.clone(), activator);

if let Some(exert_logic) = scope.worker().config().get::<trace::ExertionLogic>("differential/default_exert_logic").cloned() {
empty_trace.set_exert_logic(exert_logic);
}

let (mut reader_local, mut writer) = TraceAgent::new(empty_trace, info, logger);
let (mut reader_local, mut writer) = A::new(empty_trace, info, logger);
// Capture the reader outside the builder scope.
*reader = Some(reader_local.clone());

// Tracks the input frontier, used to populate the lower bound of new batches.
let mut prev_frontier = Antichain::from_elem(Tr::Time::minimum());
let mut prev_frontier = Antichain::from_elem(<A::Time as Timestamp>::minimum());

// For stashing input upserts, ordered increasing by time (`BinaryHeap` is a max-heap).
let mut priority_queue = BinaryHeap::<std::cmp::Reverse<(Tr::Time, K, Option<V>)>>::new();
let mut priority_queue = BinaryHeap::<std::cmp::Reverse<(A::Time, K, Option<V>)>>::new();
let mut updates = Vec::new();

move |(input, frontier), output| {
Expand Down Expand Up @@ -234,7 +255,7 @@ where
let batches = reader_local.batches_through(Antichain::new().borrow()).unwrap();
let (mut trace_cursor, trace_storage) = crate::trace::cursor::cursor_list(batches);
let mut builder = Bu::default();
let mut key_con = <BatchCursor<Tr> as Cursor>::KeyContainer::with_capacity(1);
let mut key_con = <BatchCursor<A> as Cursor>::KeyContainer::with_capacity(1);
for (key, mut list) in to_process {

key_con.clear(); key_con.push_ref(&key);
Expand All @@ -248,7 +269,7 @@ where
// Determine the prior value associated with the key.
while let Some(val) = trace_cursor.get_val(&trace_storage) {
let mut count = 0;
trace_cursor.map_times(&trace_storage, |_time, diff| count += <BatchCursor<Tr> as Cursor>::owned_diff(diff));
trace_cursor.map_times(&trace_storage, |_time, diff| count += <BatchCursor<A> as Cursor>::owned_diff(diff));
assert!(count == 0 || count == 1);
if count == 1 {
assert!(prev_value.is_none());
Expand Down Expand Up @@ -277,7 +298,7 @@ where
updates.sort();
builder.push(&mut updates);
}
let description = Description::new(prev_frontier.clone(), upper.clone(), Antichain::from_elem(Tr::Time::minimum()));
let description = Description::new(prev_frontier.clone(), upper.clone(), Antichain::from_elem(<A::Time as Timestamp>::minimum()));
let batch = crate::trace::Span::new(description, builder.done().map(Into::into));
prev_frontier.clone_from(&upper);

Expand Down
20 changes: 16 additions & 4 deletions differential-dataflow/src/operators/reduce.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ use timely::dataflow::operators::Operator;
use timely::dataflow::operators::CapabilitySet;
use timely::dataflow::channels::pact::Pipeline;

use crate::operators::arrange::{Arranged, TraceAgent};
use crate::operators::arrange::{Agent, Arranged, TraceAgent};
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,11 +75,23 @@ 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_agent::<Tr1, TraceAgent<Tr2>, T>(trace, name, tactic)
}

/// Drives a key-wise reduction like [`reduce_with_tactic`], sharing the output trace through the
/// agent `A`.
pub fn reduce_with_tactic_and_agent<'scope, Tr1, A, T>(trace: Arranged<'scope, Tr1>, name: &str, mut tactic: T) -> Arranged<'scope, A>
where
Tr1: TraceReader + 'static,
A: Agent<Time = Tr1::Time> + Clone + 'static,
A::Trace: 'static,
T: ReduceTactic<Tr1::Time, Tr1::Batch, A::Batch> + 'static,
{
let mut result_trace = None;

Expand All @@ -95,13 +107,13 @@ 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 = A::Trace::new(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);
}

let (mut output_reader, mut output_writer) = TraceAgent::new(empty, operator_info, logger);
let (mut output_reader, mut output_writer) = A::new(empty, operator_info, logger);

*result_trace = Some(output_reader.clone());

Expand Down
Loading
Loading