TimelyDataflow / TimelyDataflow/differential-dataflow

Aggregating collection into a single value

Open
#276 0 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

Dominant language
Rust
Stars
3k
Forks
211
Avg merge
10h 42m
Merged PRs (30d)
34

Description

Hi there, I've been playing around with this promising project and I'm really excited about what it can achieve.

In the examples of reduce, all aggregations seem to be per-key. I'm trying to find a way to aggregate all values in a collection to a single value. For example, summing over a collection.

For example, I've attempted this with explode, using a constant as a key:

use differential_dataflow::input::Input;
use differential_dataflow::operators::arrange::agent::TraceAgent;
use differential_dataflow::operators::arrange::{ArrangeByKey, ArrangeBySelf};
use differential_dataflow::operators::Count;
use differential_dataflow::trace::cursor::CursorDebug;
use differential_dataflow::trace::implementations::ord::OrdValSpine;
use differential_dataflow::trace::{Cursor, TraceReader};
use timely::dataflow::operators::probe::Handle;
use timely::{Configuration, PartialOrder};

pub fn sum_with_workers(initial_items: Vec<(String, isize)>) {
    timely::execute(Configuration::Process(3), move |worker| {
        let worker_index = worker.index();
        let initial_items = initial_items.clone();
        let mut probe = Handle::new();

        let (mut input_session, mut trace) = worker.dataflow(|scope| {
            let (input_session, collection) = scope.new_collection_from(initial_items);
            let trace = collection.probe_with(&mut probe).arrange_by_self().trace;

            (input_session, trace.clone())
        });

        let mut sum_trace = worker.dataflow(|scope| {
            let collection = trace
                .import(scope)
                .as_collection(|k, v| {
                    k.clone()
                })
                .probe_with(&mut probe);

            let sum = collection
                .explode(|(k, v)| {
                    Some(("CONSTANT_KEY".to_string(), v))
                })
                .count()
                .probe_with(&mut probe);

            sum.arrange_by_key().trace.clone()
        });

        let time = 2;
        input_session.advance_to(time);
        input_session.flush();

        worker.step_while(|| probe.less_than(&time));

        sum_trace.advance_by(&[time]);
        sum_trace.distinguish_since(&[time]);

        let sum_result = read_collection_at_time(&mut sum_trace, time);

        println!(
            "Worker: {}, Time: {}. Obtained sum of all items: {:?}",
            worker_index, time, sum_result
        )
    });
}

fn read_collection_at_time(
    trace_reader: &mut TraceAgent<OrdValSpine<String, isize, usize, isize>>,
    time: usize,
) -> Option<isize> {
    let (mut cursor, storage) = trace_reader.cursor();

    let mut result = None;
    while cursor.key_valid(&storage) {
        while cursor.val_valid(&storage) {
            let item = cursor.val(&storage);
            let key = cursor.key(&storage);

            let mut total = 0;
            cursor.map_times(&storage, |timestamp, update| {
                if timestamp.less_equal(&time) {
                    total = total + update;
                }
            });

            if total > 0 {
                result = Some(item);
            }

            cursor.step_val(&storage);
        }

        cursor.step_key(&storage);
    }

    result.map(|v| *v)
}

Inputting a value of

vec![
        ("id1".to_string(), 1),
        ("id2".to_string(), 2),
        ("id3".to_string(), 3),
    ];

Problem with this approach is that the result will be actual_sum * num_of_workers. I'm wondering if I should approach this differently?

Is my approach applicable for other types of aggregations that do not use explode? e.g. get the min of all items using reduce.

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start with the issue's sum_with_workers example and inspect the mentioned explode, count, and reduce entry points, focusing on how aggregation behaves across workers. Compare the constant-key approach with the requested single-value semantics. Done means a correct, worker-independent aggregation path is identified and its behavior for sum and other reductions is documented or supported.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
distributed-systems
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.