etl4s

Operators

Operator Name What it does
~> Chain a ~> b - output of a feeds into b
& / &> Fan-out a & b - run both with the same input (&> runs them concurrently)
* / *> Product a * b - run on different inputs (*> runs them concurrently)
>> Sequence a >> b - run in order, keep b's result
| Fan-in a | b - route an Either input to the matching branch
+ Choice a + b - route an Either input through independent branches
<|> Fallback a <|> b - if a throws, run b on the same input

~> chain

a b c

Output of one node feeds the next, generally the backbone of every pipeline:

val checkout = 
     parseCart ~> applyTax ~> total

& / &> fan-out (shared input)

b c (b, c)

One input, two nodes, results paired. Swap & for &> to fan out concurrently under an effect like Future. One user goes in and both cards are built from it:

val dashboard = 
     fetchUser ~> (profileCard & recentOrders) ~> layout

* / *> product (different inputs)

b d (b, d)

Where & broadcasts one input, * hands each node its own - _._1 left, _._2 right. Swap * for *> to run the two branches concurrently. Name goes to trimName, age goes to bumpAge:

val register = 
     (trimName * bumpAge) ~> saveAccount

>> sequence (keep last)

a b discarded b's result

Runs the first steps for their side effects, then flows into the real pipeline, keeping its result. This runs the setups steps in order then flows into the real pipeline:

val ingest = 
     clearStaging >> warmCache >> (extract ~> transform ~> load)

| fan-in (route an Either in)

b Left c Right out

An Either input goes left or right, both branches merging to one type. A legacy id or a uuid comes in, one user comes out:

val load = 
     (byLegacyId | byUuid) ~> loadAccount

+ choice

b Left c Right Left' Right'

Routes an Either through independent branches and keeps the Either on the way out. Each txn kind is handled on its own branch, then merged back with |

val settle = 
     classifyTxn ~> (handleRefund + handleCharge) ~> (fileRefund | postCharge)

<|> fallback (try, then recover)

a b on throw result

Runs the left node; if it throws, runs the right on the original input. This tries fetchLive, and falls back to fetchCached if it throws:

val convert = 
     (fetchLive <|> fetchCached) ~> applyRates ~> total

Under an error-tracking effect the recovery flows through that effect's error channel instead of a thrown exception.

Concurrency needs a concurrent effect

&> and *> only run their branches concurrently when compiled to a concurrent effect such as Future (or IO), via .compile[Future]. Plain unsafeRun (the Id interpreter) has no threads, so it runs them in sequence.

Auto-flatten and .zip

Chaining fan-outs auto-flattens the tuple: a & b & c has type Node[X, (A, B, C)] (not ((A, B), C)), and the same holds for &>:

import etl4s._

val n1 = Node[String, Int](_.length)
val n2 = Node[String, String](_.toUpperCase)
val n3 = Node[String, Boolean](_.nonEmpty)

val flat = n1 & n2 & n3
flat.unsafeRun("hi")
flat has type Node[String, (Int, String, Boolean)], and you will get:
(2, "HI", true)
If you already have a node whose output is a nested tuple, .zip flattens it:
val nested  = (n1 & n2) & n3
val flatten = nested.zip
nested has type Node[String, ((Int, String), Boolean)] and flatten has type Node[String, (Int, String, Boolean)].