Skip to content

Pipeline execution order

vird edited this page Aug 6, 2017 · 12 revisions

NOTE This is only theory right now. Full implementation will be later.

Sync

1-frame

Simplest example

["hello world"] | stdout

This code will have coffee-script equivalent

for v in ["hello world"]
  console.log v

and can be optimized to

console.log "hello world"

because we detected that frame count == 1

Alias

You can use identifier for sources, sinks and intermediate connections. There are equivalent pipelines:

src = ["hello world"]
src | stdout

Same

["hello world"] | dst
dst | stdout

Split

You can make some calculation once and reuse it with alias

gramm | (t)->t/1000 | kilogramm
gramm | (t)->t/100000 | ton

You can code less with joined pipes. Here is split example

gramm | (t)->t/1000 | kilogramm
      | (t)->t/100000 | ton

This can make your code 1 useless variable less.

Join

You can use 2 result from different alias

1 | a
2 | b
a+b | c

You can use join syntax

1 |
2 |
  | $1+$2 | c

This combos with split syntax

src | make_op1 |
    | make_op2 |
               | $1+$2 | dst

So you can make more visual-friendy diagramms. This code will have coffee-script equivalent

dst = []
for v in src
  tmp1 = make_op1 v
  tmp2 = make_op2 v
  dst.push tmp1+tmp2
  reduce_accumulator1 += v

BEWARE | symbol must be at same position. Tabs are not allowed. In strict mode | allowed only at %10==0 columns for forcing better formating.

Note. Join depends on which line you make continuation.

  • Separate line.
    • By default checks that all branches has same frame count.
    • There are some special join functions that can make specific operations over results in different branches.
      • join
      • cross join
  • Shared line with some branch. Will take frame count from this line. All other lines accessible with aliasing.

Clocking

Clocking is important for understanding how pipeline will be unrolled into sequentian code. Let's try multiframe example

[1,2,3] | stdout

This code will have coffee-script equivalent

for v in [1,2,3]
  console.log v

We can add multiple functions. All functions by default accept 1 frame and produce 1 frame, so count frame in == count frame out.

[1,2,3] | (t)->t+1 | stdout

This code will have coffee-script equivalent

for v in [1,2,3]
  res2 = v+1
  console.log res2

Nothing special we just applied (t)->t+1 to all elements in array (in pipeline terms to all frames in stream). Let's meet reduce

[1,2,3] | reduce0 0, (a,b)->a+b | stdout
[1,2,3] | reduce0 0, (+) | stdout # same with syntax sugar

This code will have coffee-script equivalent

reduce_accumulator1 = 0
for v in [1,2,3]
  reduce_accumulator1 += v
console.log reduce_accumulator1

Note that some code generated before main loop and some after. That's because initial pipeline has different frame count on different stages.

[1,2,3] | reduce0 0, (a,b)->a+b | stdout
# 3     | 3                   1 | 1

Things became more complex with split/join.

[1,2,3] | reduce0 0, (+) | avg | 
        |                |     | (t)->t-avg | dst

This code will have coffee-script equivalent

reduce_accumulator1 = 0
src = [1,2,3]
for v in src
  reduce_accumulator1 += v
avg = reduce_accumulator1
dst = []
for v in src
  dst.push v-avg

Note that (t)->t-avg has exactly same frame count as [1,2,3], but other loop created. Join will check that all branches have same frame count. And only at this case it produce 1 loop. All other cases there would be multiple loops

Async

Any pipeline block produces a process that can be assigned to a variable. You can operate process with API.

proc = a | b | c
proc.start()
proc.pause()
proc.stop()

If pipeline is not assigned to any variable, it starts implicitly.\ Async makes this more interesting, but have pretty same syntax with a lot different code generated.

each 1h | "hello world" | stdout

Here we used new things. each is async function that use setInterval and emit frame each 1 hour. "hello world" const replaces current frame with const. CoffeeScript equivalent

setTimeout ()->
  console.log "hello world"
, 60*60*1000

using sink alias in non-pipeline code will force barrier creation. So you will wait until sink would be resolved

[1,2,3] | rate 1rps | dst
# some code
console.log dst

in CoffeeScript

src = [1,2,3]
anonymous_process = (on_end)->
  res = []
  for v in src
    dst.push v
    await setTimeout defer(), 1000
  on_end null res
# some code
await anonymous_process defer(err, dst); throw err if err
console.log dst

You can make lazy calculations with pipelines. Alias with no use will be cutted by compiler. Alias that are not 100% be accessed will call function-wrapper that have all needed logic only if that calculation is really needed.

Experimental

argument count + type autodetect

[1,2,3] | (t)->t            | res # this is map
[1,2,3] | (a,b)->a,b        | res # this is reduce
[1,2,3] | (a,cb)->cb(a)     | res # this is async map
[1,2,3] | (a,b,cb)->cb(a+b) | res # this is async reduce

Clone this wiki locally