etl4s
API Reference
  • .If(pred)(branch)
  • .ElseIf(pred)(branch)
  • .Else(branch)
  • If[A](pred)(branch)
  • IfCtx[C](pred)(branch)
  • .IfCtx(pred)(branch)
  • .ElseIfCtx(pred)(branch)

Conditional branching

Route data down different pipelines with If, ElseIf, and Else. Branch on the data flowing through, or on outside config.

Branch on data

The predicate runs against the input, and only the matching branch fires:

val classify = score
  .If(_ > 0)     (tagPositive)
  .ElseIf(_ < 0) (tagNegative)
  .Else          (tagZero)

Close with .Else and it is exhaustive. Drop the .Else and unmatched input flows straight through:

val maybeBoost = enrich.If(_.score > 10)(applyBoost)

Start with a branch

If[A] opens a pipeline on the branch itself, no upstream node needed:

val ship = If[Order](_.isRush)(expedite).Else(standard) ~> notify

Branches are pipelines

Any branch can be a whole pipeline, fan-out and all:

val offers = extractUser
  .If(_.tier == "premium")      (validate ~> enrich ~> premiumOffer)
  .ElseIf(_.tier == "standard") (validate ~> standardOffer)
  .Else                         (guestNotice)
val profile = loadUser
  .If(_.wantsDetails) ((self & loadMetrics & loadHistory) ~> fullProfile)
  .Else               (simpleProfile)

Swap & for &> to fan out concurrently under an effect.

Branch on config and data

When the decision needs config too, take a typed condition (cfg: Config) => (data: A) => Boolean:

val route = source
  .If((cfg: Config) => (n: Int) => n < cfg.threshold) (formatBelow)
  .Else                                               (formatAbove)

route.provide(Config(10)).unsafeRun(5) // "below:5"

Branch on config alone

When only the config matters and the data is irrelevant, use IfCtx / ElseIfCtx - the condition is just Config => Boolean. IfCtx[Config] starts the pipeline on the context itself, no source node:

val ingest =
  IfCtx[Config](_.isBackfill)(readSnapshot ~> replay ~> load)
    .ElseIfCtx(_.isDryRun)   (readStream ~> validate ~> logOnly)
    .Else                    (readStream ~> validate ~> load)

ingest.provide(Config(isBackfill = true, isDryRun = false)).unsafeRun(batch)

Already have an upstream Reader? Call .IfCtx on it instead:

val ingest = source
  .IfCtx(_.isBackfill)  (readSnapshot ~> replay ~> load)
  .ElseIfCtx(_.isDryRun)(readStream ~> validate ~> logOnly)
  .Else                 (readStream ~> validate ~> load)

Scala 2 vs Scala 3

The API is identical across versions, but Scala 3's type system enables more flexibility.

Scala 3: branches can return different types (union):

val router = score
  .If(_ > 0)     (toLabel)    // String
  .ElseIf(_ < 0) (negate)     // Int
  .Else          (toDouble)   // Double

The result type is Node[Int, String | Int | Double].

Scala 3: config-aware branches accumulate their config via intersection (&): mixing branches that need DbConfig and CacheConfig yields a pipeline that must be provided DbConfig & CacheConfig.

Scala 2

All branches must return the same type, and share the same config type:

val router = score
  .If(_ > 0)     (posLabel)
  .ElseIf(_ < 0) (negLabel)
  .Else          (zeroLabel)