-
.ensure(input, output, change) -
.ensurePar(input, output) -
ValidationException
Ensurers
When writing dataflows, you often want to validate inputs and outputs at runtime - and reuse those validations across nodes, collecting all errors instead of failing on the first.
Validators are just functions A => Option[String]. Return None if valid, Some("error message") if not:
import etl4s._
val isPositive = (x: Int) => if (x > 0) None else Some("Must be positive")
val lessThan1k = (x: Int) => if (x < 1000) None else Some("Must be < 1000")
val notEmpty = (s: String) => if (s.nonEmpty) None else Some("Cannot be empty")
.ensure() lets you attach validators to any Node:
val process = Node[Int, String](n => s"Value: $n")
.ensure(
input = Seq(isPositive, lessThan1k),
output = Seq(notEmpty)
)
process.unsafeRun(42)
process.unsafeRun(-5)
Change Validation
Validate by examining both input and output together. The change validator receives a tuple (input, output):
/* Ensure deduplication never grows the list */
val noGrowth: ((List[Int], List[Int])) => Option[String] = {
case (in, out) =>
if (out.size <= in.size) None
else Some(s"Output grew: ${in.size} -> ${out.size}")
}
val dedupe = Node[List[Int], List[Int]](_.distinct)
.ensure(change = Seq(noGrowth))
dedupe.unsafeRun(List(1, 2, 2, 3))
Error Accumulation
Multiple failures are collected:
val lessThan100 = (x: Int) => if (x < 100) None else Some("Must be < 100")
val isEven = (x: Int) => if (x % 2 == 0) None else Some("Must be even")
val validate = Node[Int, Int](identity)
.ensure(input = Seq(isPositive, lessThan100, isEven))
validate.unsafeRun(-5)
You will get:
Validation under effects
The default .unsafeRun throws on failure. To capture the failure as a value,
run through an effect. Under .compile[Try] a validation failure surfaces as a
Failure(ValidationException):
import etl4s._
import scala.util.{Try, Failure}
val node = Node[Int, String](_.toString)
.ensure(input = Seq(isPositive))
node.compile[Try].unsafeRun(42)
node.compile[Try].unsafeRun(-5)
See Effect polymorphism for the full list of effects.
Parallel Validation
Use .ensurePar() in place of .ensure() to mark the checks within each stage
as eligible to run concurrently. This only actually runs them in parallel under a
concurrent effect (e.g. .compile[Future] etc).
import etl4s._
import scala.concurrent.{Future, Await}
import scala.concurrent.duration._
import scala.concurrent.ExecutionContext.Implicits.global
/* Two independent, potentially expensive checks */
val isPositive = (x: Int) => if (x > 0) None else Some("Must be positive")
val lessThan100 = (x: Int) => if (x < 100) None else Some("Must be < 100")
val validate = Node[Int, Int](identity)
.ensurePar(input = Seq(isPositive, lessThan100))
Await.result(validate.compile[Future].unsafeRun(42), 5.seconds)
You will get:
Handling Failures
Validation failures throw a ValidationException. Recover with .onFailure()
or wrap the run in your own Try:
import scala.util.Try
val node = Node[Int, String](_.toString)
.ensure(input = Seq(isPositive))
Try(node.unsafeRun(-5))
You will get:
Config-Aware Validation
Ensurers work on config nodes too. Validators are curried Config => A => Option[String] so they can access config:
case class Config(minValue: Int, maxValue: Int)
val inRange: Config => Int => Option[String] = cfg => n =>
if (n >= cfg.minValue && n <= cfg.maxValue) None
else Some(s"Must be between ${cfg.minValue} and ${cfg.maxValue}")
val process = Transform[Int, Int].requires[Config] { cfg => n => n * 2 }
.ensure(input = Seq(inRange))
process.provide(Config(0, 100)).unsafeRun(50)
process.provide(Config(0, 100)).unsafeRun(150)
etl4s