std/parallel/lib
std/parallel/src/lib.trb
One thing going faster: Parallel, a pipeline whose stages run on the workers of the pool, the parallel() that
starts one from anything that can be iterated, and Cut, what a collection implements to be cut into chunks without
being copied first.
const squares = (0..1000).parallel().map({ _ * _ }).toList().await()
print squares.length()
The words are the ones of the sequential pipeline - map, filter, filterMap, collect, toList, count, sum,
minBy, maxBy, find, forEach - and a terminal answers a Task, because the work runs while the caller's worker
is free to do something else. Results arrive in input order, always: the input is cut into chunks whose borders depend
on its length alone (at most 64 of them), every chunk runs as a task an idle worker may take, and the chunks are read
back in order. So the answer never depends on the number of workers, the load of the machine or which chunk finished
first (docs/design/CONCURRENCY.md section 4).
Where the chunks come from. A Cut cuts itself: a List and an Array into lists of their items, a Range of
integers arithmetically into smaller ranges - (0..1000000).parallel() never builds a list of a million numbers.
Everything else that can be iterated - a Set, a Map, the characters of a String, a pipeline, a type of the
program's own - is read once into a list, which is then cut the same way. Either way the borders are the same function
of the length, so the answer is too.
One pass per chunk. The stages of a pipeline are fused: map followed by filter runs every chunk once, item by
item through both closures, and builds no list between them. What a chunk hands back is what the terminal keeps of
it - its items for toList, a number for count and sum, at most one item for find, minBy and maxBy.
Pitfalls
- A chunk moves to another worker only where its piece and the closures of the pipeline may cross one: plain items
(
Int,Float, a range, a record of those) as they are,Strings, lists, variants and records of those as a copy made when the chunk's task starts, and so is the environment of a closure that captures them. A pipeline overshared typeobjects, or with a closure that captures one, avaror a trait-typed value, runs its chunks on the caller's worker, one after the other - the same answer, without the speed-up (docs/design/CONCURRENCY.md section 16, "The copy at the crossing"). collectruns the stages on the workers and the accumulator on the caller's worker, one chunk after the other, joining the chunks' outputs with the accumulator'smerge: an accumulator is a trait-typed value, which cannot cross to another worker yet.- Cancelling the task a terminal answers stops every chunk at its next suspension point or loop turn.
Related
Workers.count- how many workers the pool has, which is whatworkers:defaults to.Merge- whatcollectneeds of an accumulator.spawn- one task, whereparallel()is many of the same.
trait Cut
trait Cut<Item, Piece: Iterate<Item>> with Iterate<Item>, Length
A collection that cuts itself into pieces without being read first: its length, and the piece of it the positions of
a range cover. parallel() cuts one into chunks by its length alone; anything else is read into a list first.
Piece is whatever carries a part of the collection best - a List and an Array answer an ArrayList of the items,
which a worker may take as it is, and a Range of integers answers a smaller range, so cutting it costs two additions.
What a chunk crosses to another worker as is the piece, which is why it is a type of its own rather than Self.
Examples
type Stretch with Cut<Int, Range<Int>> {
first: Int
count: Int
fn iterate(): Iterator<Int> {
(first..(first + count)).iterate()
}
fn length(): Int {
count
}
fn cut(range: Range<Int>): Range<Int> {
(first + range.start)..(first + range.end)
}
}
print Stretch(5, 10).parallel().sum().await()
Pitfalls
- The piece has to hold the same items, in the same order, as iterating the collection over those positions: the
answer of
parallel()is the sequential one only if cutting and iterating agree.
fn cut
fn cut(range: Range<Int>): Piece
The items at the positions range covers, as a piece of their own. range lies within 0..length().
extend List<Item> with Cut<Item, ArrayList<Item>>
extend<Item> List<Item> with Cut<Item, ArrayList<Item>>
A list cuts into lists of its items: a piece holds nothing the list holds, so a worker may take it as it is.
fn cut
fn cut(range: Range<Int>): ArrayList<Item>
extend Array<Item, Size> with Cut<Item, ArrayList<Item>>
extend<Item, const Size: Int> Array<Item, Size> with Cut<Item, ArrayList<Item>>
An array cuts into lists of its items, the way a list does.
fn cut
fn cut(range: Range<Int>): ArrayList<Item>
extend Range<Int> with Cut<Int, Range<Int>>
extend Range<Int> with Cut<Int, Range<Int>>
A range of integers cuts into smaller ranges: position index is the value start + index.
fn cut
fn cut(range: Range<Int>): Range<Int>
extend Iterate<Item>
extend<Item> Iterate<Item>
Spreads a pipeline over the workers: anything that can be iterated.
fn parallel
fn parallel(workers: Int = Workers.count(), chunk: Int? = None): Parallel<Item>
This pipeline, spread over the workers: its items are read into a list once, which is cut into chunk pieces (at
most 64 where it is None, and never more than there are items), and at most workers of them run at the same
time. The borders depend on the number of items and on chunk alone, never on workers, so neither does the answer.
Examples
const letters: Set<String> = ["a", "bb", "ccc"]
print letters.parallel(workers: 2).map({ _.byteLength() }).sum().await()
extend Cut<Item, Piece>
extend<Item, Piece: Iterate<Item>> Cut<Item, Piece>
Spreads a collection that cuts itself over the workers, without reading it first.
fn parallel
fn parallel(workers: Int = Workers.count(), chunk: Int? = None): Parallel<Item>
This collection, spread over the workers: it is cut into chunk pieces (at most 64 where it is None, and never
more than it has items), and at most workers of them run at the same time. The borders depend on length() and
chunk alone, never on workers, so neither does the answer.
Examples
const total = (1..=100).parallel(workers: 2).map({ _ * 2 }).sum().await()
print total
extend List<Item>
extend<Item> List<Item>
Spreads a list over the workers. A list is a Cut, and this is that parallel() again: a value of the List trait
cannot be handed over as a value of Cut - Cut is implemented for List and is no supertrait of it, so no table a
list carries holds it - so it calls its own cut instead.
fn parallel
fn parallel(workers: Int = Workers.count(), chunk: Int? = None): Parallel<Item>
This list, spread over the workers: the items are cut into chunk pieces (at most 64 where it is None, and never
more than there are items), and at most workers of them run at the same time. The borders depend on the number of
items and on chunk alone, never on workers, so neither does the answer.
Examples
const lengths = ["a", "bb", "ccc"].parallel(workers: 2).map({ _.byteLength() }).toList().await()
print lengths
type Parallel
type Parallel<Item>
A pipeline whose stages run on several workers, with the results in input order. Made by parallel() on anything
that can be iterated; nothing runs until a terminal - toList, collect, count, sum, minBy, maxBy, find,
forEach - is called, and the terminal answers a Task. The stages before it run in one pass per chunk.
Examples
const total = (1..=100).parallel().map({ _ * 2 }).sum().await()
print total
fn map
fn map<Output>(transform: Transform<Item, Output>): Parallel<Output>
Every item through transform, the chunks on the workers.
fn filter
fn filter(predicate: Predicate<Item>): Parallel<Item>
The items for which predicate answers true, in their order.
fn filterMap
fn filterMap<Output>(transform: Transform<Item, Output?>): Parallel<Output>
Every item through transform, keeping the values that are there.
fn toList
fn toList(): Task<List<Item>>
Every item, in input order.
fn collect
fn collect<Output>(collector: Merge<Item, Output>): Task<Output>
The items pushed into collector: a copy of it per chunk that received an item, in input order, and the chunks'
outputs joined left to right with its merge. A pipeline without any item answers the collector's own finish(),
and merge never sees the output of a chunk the stages left empty (the contract of Merge). The stages run on the
workers; the collector runs on the caller's worker.
Examples
const words = ["pear", "fig", "plum"]
print words.parallel().map({ _.byteLength() }).collect(counting()).await()
fn count
fn count(): Task<Int>
How many items come out of the pipeline.
fn sum
fn sum(): Task<Item> where Item: Add & From<Int>
The items added up: each chunk on its worker, then the chunks' sums in input order. A Float sum is therefore the
same number on every machine, and not the same number as the sequential sum(), which adds in another order.
fn minBy
fn minBy<Key: Compare>(key: Transform<Item, Key>): Task<Item?>
The item with the smallest key, or None. Of two equal keys the earlier one in input order, as minBy answers.
fn maxBy
fn maxBy<Key: Compare>(key: Transform<Item, Key>): Task<Item?>
The item with the largest key, or None. Of two equal keys the earlier one in input order, as maxBy answers.
fn find
fn find(predicate: Predicate<Item>): Task<Item?>
The first item in input order for which predicate answers true, or None. A chunk stops at its first match, and
no further chunk is started once every chunk before it has answered and one of them found one.
fn forEach
fn forEach(body: Action<Item>): Task<Void>
Runs body for every item, the chunks on the workers. The order in which items of two chunks run is the machine's.