std/iteration/stage
std/iteration/src/stage.trb
A stage is the middle of a pipeline, written once for both worlds: it turns an Accumulator<Output, Final> into
an Accumulator<Input, Final>, a transducer in Clojure's word.
const interesting = filtering<Event> { _.level >= .Warning }.then(mapping { _.message })
const fromList = events.through(interesting).toList() // a list, synchronously
const fromBody = response.body.through(interesting).toList() // an HTTP body, asynchronously
Because a stage never asks where its values come from, every stage exists once: the pull side (Iterate) and the
push side (Source) only differ in the two drivers that feed one. And because a stage is an ordinary value, a
pipeline is a value: it can be named, passed, stored in a field and applied more than once.
Open
The per-stage iterators of stages.trb (Mapped, Filtered, ...) still carry Iterate.map and friends of their
own, next to the Stage values here that do the same job through through(...). The two live side by side and
mean the same thing.
Related
Accumulator- what a stage wraps and hands values to;ontois how a stage sits in front of one.
trait Stage
trait Stage<Input, Output>
onto is the whole protocol, and it is the reason a stage composes: what it answers is an Accumulator again, so
the next stage can wrap that one. Nothing here is asynchronous - a stage computes, it never waits.
Examples
const evens: Stage<Int, Int> = filtering { _ % 2 == 0 }
const numbers = [1, 2, 3, 4, 5, 6]
print numbers.through(evens).toList()
fn onto
fn onto<Final>(downstream: Accumulator<Output, Final>): Accumulator<Input, Final>
Wraps downstream, so that every Input pushed into the answer arrives as zero, one or many Outputs in
downstream. The wrapper has to pass finish() on exactly once, and isDone() from downstream upwards, so
that taking and first end the pipeline instead of pulling into nothing.
fn then
fn then<Final>(other: Stage<Output, Final>): Stage<Input, Final>
filtering { ... }.then(mapping { ... }): two stages as one, which is what makes a pipeline a value.
fn mapping
fn mapping<Input, Output>(transform: Transform<Input, Output>): Stage<Input, Output>
Every stage is a function, not a member of Stage, for one reason: a member of a trait used as a namespace has to
fix the trait's own type arguments, and Stage.filter could not (its answer says nothing about Output). The names
are the gerunds the collectors already use (counting, joining, groupingBy), which also keeps them apart from
the methods of the same meaning on Iterate and Source (items.map(f) is items.through(mapping(f))).
fn filtering
fn filtering<Item>(predicate: Predicate<Item>): Stage<Item, Item>
Keeps every value that answers the predicate, and drops the rest.
fn filterMapping
fn filterMapping<Input, Output>(transform: Transform<Input, Output?>): Stage<Input, Output>
map and filter in one: keeps the value of every Some. The bridge for functions that answer an Option.
fn mappingWhile
fn mappingWhile<Input, Output>(transform: Transform<Input, Output?>): Stage<Input, Output>
Like filterMapping, but the first None ends the pipeline.
fn flatMapping
fn flatMapping<Input, Output>(transform: Transform<Input, Iterate<Output>>): Stage<Input, Output>
Transforms every value into its own Iterate, and flattens all of them into one sequence.
fn taking
fn taking<Item>(amount: Int): Stage<Item, Item>
Ends the pipeline after amount values, and asks for none of the values behind them.
fn takingWhile
fn takingWhile<Item>(predicate: Predicate<Item>): Stage<Item, Item>
Values up to, but not including, the first one that fails the predicate.
fn skipping
fn skipping<Item>(amount: Int): Stage<Item, Item>
Every value after the first amount, which are read and thrown away.
fn indexing
fn indexing<Item>(): Stage<Item, (index: Int, item: Item)>
(0, first), (1, second), ... with the halves named index and item - the stage behind indexed().
fn chunking
fn chunking<Item>(size: Int): Stage<Item, List<Item>>
Groups values into lists of at most size. The last group is whatever is left, so nothing is dropped, and a
size below one is a group per value.