Advanced
Pipelines are values
Pipelines being values unlocks some powerful niceties. For example given:
Unit test pipeline shape
test("etl graph is wired as designed") {
assertEquals(billing.stages.map(_.name), List("parse", "applyTax", "format"))
assertEquals(billing.stages.map(s => s.in -> s.out).last, "Double" -> "String")
}
Govern dataflow architecture
def audit(p: Node[?, ?]): Unit = {
val forbidden = p.stages.filter(_.fullName.startsWith("com.acme.legacy"))
require(forbidden.isEmpty, s"pipeline pulls in banned stages: ${forbidden.map(_.name)}")
}
Generate docs that never drift
Higher Order Nodes
import etl4s._
object CustomerOps {
def activeOnly =
Node[List[Customer], List[Customer]](_.filter(_.isActive))
def topSpenders(n: Int) =
Node[List[Customer], List[Customer]](_.sortBy(-_.spend).take(n))
def inRegion(region: String) =
Node[List[Customer], List[Customer]](_.filter(_.region == region))
}
import CustomerOps._
val pipeline =
extract ~> activeOnly ~> inRegion("EU") ~> topSpenders(100) ~> load
Dynamic Composition
You can build pipelines at runtime instead of writing every ~> by hand.
import etl4s._
val rules: List[Node[Row, Row]] = List(
trimStrings,
dropEmpty,
normalizeDates
)
val cleaningRules: Node[Row, Row] = rules.reduce(_ ~> _)
val pipeline =
extract ~> cleaningRules ~> load
Since reduce throws on empty lists, you can fold from Node.identity (no-op Node)
and the result is a valid pipeline just with zero steps
Assemble custom pipelines based on some configuration type
case class Config(dedupe: Boolean, enrich: Boolean)
def buildPipeline(cfg: Config): Node[Row, Row] = {
val optional = List(
cfg.dedupe -> dedupe,
cfg.enrich -> enrich
)
optional
.collect { case (true, step) => step }
.foldLeft(Node.identity[Row])(_ ~> _)
}
Custom Operators
import etl4s._
extension [A, B](node: Node[A, B]) {
def timed(label: String): Node[A, B] = Node { input =>
val start = System.currentTimeMillis()
val result = node(input)
println(s"$label: ${System.currentTimeMillis() - start}ms")
result
}
}
val pipeline =
extract ~> transform.timed("main") ~> load
Symbolic Operators
Define your own symbolic operators like !! and @@ below
etl4s