Parallel Execution
etl4s has an elegant shorthand for grouping and parallelizing operations that share the same input type:
import etl4s._
/* Simulate slow IO operations (e.g: DB calls, API requests) */
val e1 = Extract { Thread.sleep(100); 42 }
val e2 = Extract { Thread.sleep(100); "hello" }
val e3 = Extract { Thread.sleep(100); true }
Sequential run of e1, e2, and e3 (~300ms total):
Parallel run of e1, e2, e3 on their own JVM threads with Scala Futures
(~100ms total, same result, 3X faster):
Mix sequential and parallel execution - first two parallel (~100ms), then the third (~100ms):
Full example of a parallel pipeline:
val consoleLoad: Load[String, Unit] = Load(println(_))
val dbLoad: Load[String, Unit] = Load(x => println(s"DB Load: ${x}"))
val merge = Transform[(Int, String, Boolean), String] { case (i, s, b) =>
s"$i-$s-$b"
}
val pipeline =
(e1 &> e2 &> e3) ~> merge ~> (consoleLoad &> dbLoad)
pipeline.compile[Future].unsafeRun(())
Concurrency only kicks in under a concurrent effect - see
Effect polymorphism. For per-element parallelism over a
collection, see eachPar in Batch collections.
etl4s