Reference

std/iteration/collectors

std/iteration/src/collectors.trb

Collectors decide where the values of a pipeline end up:

employees.collect(groupingBy { _.department }.then(averaging { _.salary }))
words.collect(joining(", "))

An Accumulator is both: the description of a run and the state of one. Values are pushed into it one at a time through add, and finish() produces the result. Because it is a value, a copy of it is a run of its own - which is why collect may be called twice with the same accumulator and why groupingBy(...).then(downstream) copies the downstream per group. Because values are pushed, an Accumulator is the synchronous counterpart of Sink, which pushes the same way but may have to wait for the other end to keep up.

A collection is not one, and there is no trait for "something with add" either: an accumulator is a type of its own beside the collections, the way Collector/Collectors.toList() are in Java. ListAccumulator is the one this package ships; a package that owns a collection may ship its own beside it, and a caller names it.

Examples

const words = ["pear", "kiwi", "plum", "fig"]
print words.collect(groupingBy { _.byteLength() })

trait Accumulator

trait Accumulator<Item, Output>

What to do with the values of a pipeline, and the state of doing it: values arrive one at a time through add, and finish produces the result once they have all arrived. collect(accumulator) is the general terminal operation of Iterate - toList, count, fold and the rest are all written on top of it.

A description and a run are one type, because a value is a copy. collect fills a copy of what it was given, so one accumulator may be collected with twice and the two runs never see each other, and then(downstream) needs no factory: it keeps the downstream and copies it per group.

It stands alone. A collection does not implement it: a container that is merely filled has no result of a run to give and no isDone() to answer, so what gathers into one is a type beside it (ListAccumulator, into).

Examples

const numbers = [1, 2, 3, 4, 5]
print numbers.collect(summing { value: Int => value })

Pitfalls

  • A description and a half-filled run have the same type, so an accumulator that has already been added to and is then handed to collect or to then starts every run from what is in it. Build one where it is used (collect(listing())), the way every collector of this package answers a fresh value.

Related

  • ListAccumulator - the one this package ships, and the shape a package writes for its own collection.
  • groupingBy - the one accumulator with a then(downstream) for a different one per group.

fn add

var fn add(value: Item)

Pushes one more value in.

fn finish

fn finish(): Output

The result, once every value has been pushed in. Called exactly once, after the last add.

fn isDone

fn isDone(): Bool

Whether more values would change the result. A driver asks before it pulls the first value and after every add, and stops as soon as the answer is true - which is what lets taking, first and find end a pipeline over an infinite or expensive source instead of reading to the end and throwing the rest away.

The default is false: an accumulator that sums, counts or collects is never done. A wrapper (std/iteration's stages) answers what the accumulator below it answers, plus whatever it knows itself.

trait Merge

trait Merge<Item, Output> with Accumulator<Item, Output>

Two partial results of one accumulator, joined: what lets an accumulator run over pieces of its input - the chunks of parallel(), a divide-and-conquer fold, a map-reduce over the parts of a file - and still answer what one run over the whole input answers.

The contract is three lines. merge is associative: merge(merge(a, b), c) is merge(a, merge(b, c)). It need not be commutative, because it is called in the order of the pieces, left to right. And it is never called with the output of a piece that received no value: such a piece is skipped, and an input without any value takes the accumulator's own finish() instead. That third line is what makes joining(separator:, prefix:, suffix:) mergeable at all - without it every merge would have to decide whether a separator belongs between two pieces one of which is not there.

It joins two outputs and not two accumulators: two trait-typed accumulators need not have the same type, and the output is what comes back from wherever a piece ran anyway. A collector whose output lost what a merge would need - averaging, whose count is gone once it divided - has none; the shape of the workaround is a mergingCollector whose output keeps it.

Examples

const counter = counting<Int>()
print counter.merge([1, 2].collect(counter), [3].collect(counter))

Related

fn merge

fn merge(first: Output, second: Output): Output

The output of the pieces before and the output of the piece after them, joined.

fn into

fn into<Target: From<Iterate<Item>>, Item>(): Accumulator<Item, Target>

Collects into anything a pipeline can be turned into (From<Iterate<Item>>), which every collection is. Target can be given on its own, the rest is inferred.

It gathers into a List and calls Target.from once at the end, which is the answer that works for every target. A package that owns a collection may ship an accumulator that writes straight into it - ListAccumulator is the shape - and a caller who wants that one names it (collect(ListAccumulator<Int>())). Nothing picks it up automatically: which run a pipeline makes is written at the call, not inferred from the target.

Examples

const numbers = [3, 1, 2, 1]
print numbers.collect(into<Set<Int>>())

fn listing

fn listing<Item>(): Merge<Item, List<Item>>

Shorthand for ListAccumulator<Item>(), which is what a pipeline ends in most often. Two lists merge by concatenation.

type ListAccumulator

type ListAccumulator<Item> with Merge<Item, List<Item>>

Gathers the values into a List and answers it: the specialised accumulator this package ships, and the shape a package that owns a collection writes for its own (HashSetAccumulator writing straight into a set).

Examples

print([3, 1, 2].collect(ListAccumulator<Int>()))

fn add

var fn add(value: Item)

Appends to the list that is being gathered.

fn finish

fn finish(): List<Item>

The list, as it stands after the last add.

fn merge

fn merge(first: List<Item>, second: List<Item>): List<Item>

The two lists one after the other.

fn collector

fn collector<Item, State, Output>(initial: State, finish: (State) => Output, step: (State, Item) => State): Accumulator<Item, Output>

State is the starting point of every run. Every run works on its own copy.

fn counting<Item>(): Accumulator<Item, Int> {
  collector(0, finish: { _ }) { count, _ => count + 1 }
}

fn mergingCollector

fn mergingCollector<Item, State, Output>(initial: State, finish: (State) => Output, merge: (Output, Output) => Output, step: (State, Item) => State): Merge<Item, Output>

collector with a Merge: merge joins the outputs of two pieces of the input, in their order, and is never handed the output of a piece that received no value. It has to be associative, so that the answer does not depend on where the pieces were cut.

An average merges once its output keeps the count:

const averaged = mergingCollector((0.0, 0), finish: { _ }, merge: { first, second =>
  (first.0 + second.0, first.1 + second.1)
}) { state, value: Float => (state.0 + value, state.1 + 1) }
const total = [1.0, 2.0, 6.0].collect(averaged)
print(total.0 / Float.from(total.1))

fn counting

fn counting<Item>(): Merge<Item, Int>

How many values there are. Two counts merge by +.

fn summing

fn summing<Item, Total: Add & From<Int>>(value: Transform<Item, Total>): Merge<Item, Total>

Every value, added together, starting from zero. Two sums merge by +.

fn averaging

fn averaging<Item>(value: Transform<Item, Float>): Accumulator<Item, Float?>

None for an empty pipeline.

fn minBy

fn minBy<Item, Key: Compare>(key: Transform<Item, Key>): Merge<Item, Item?>

The value with the smallest key, or None for an empty pipeline. Keeps the earlier of two equal keys, and so does its merge: the second piece's value wins only with a smaller key.

fn maxBy

fn maxBy<Item, Key: Compare>(key: Transform<Item, Key>): Merge<Item, Item?>

The value with the largest key, or None for an empty pipeline. Keeps the earlier of two equal keys, and so does its merge: the second piece's value wins only with a larger key.

fn joining

fn joining(separator: String = "", prefix: String = "", suffix: String = ""): Merge<String, String>

Every value, joined into one String, with separator between values and prefix/suffix around all of them.

It buffers the pieces and concatenates once at the end, for the reason concatenated gives: a String is a value, so appending to one copies everything that is already in it.

fn partitioningBy

fn partitioningBy<Item>(predicate: Predicate<Item>): Merge<Item, (List<Item>, List<Item>)>

(matching, rest). Two pairs merge list by list.

fn groupingBy

fn groupingBy<Item, Key: Hash>(key: Transform<Item, Key>): Grouping<Item, Key>
groupingBy { _.department }                              // Map<String, List<Employee>>
groupingBy { _.department }.then(counting())             // Map<String, Int>
groupingBy { _.department }.then(maxBy { _.salary }) // Map<String, Employee?>