Reference

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

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.