Skip to content

Latest commit

 

History

8 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Free Monads for Declarative Data Pipelines

Overview

The problem

We want to write a pipeline like this:

for
  records <- readFromKafka(topic)
  cleaned <- transformData(records)
  _       <- writeToDatabase(cleaned)
yield ()

We also want two more things:

  1. We write the pipeline once and run it in more than one way. Real Kafka in production. An in-memory list in tests. Print only for a dry run.
  2. We inspect the pipeline before it runs. We log it, count the steps, and check it.

Normal code cannot do this. When we call readFromKafka, it reads from Kafka at once. There is nothing to inspect and nothing to replace.

The solution

We do not run the instructions. We record them. The pipeline becomes a description. An interpreter reads the description and does the work. We can replace the interpreter and keep the pipeline.

The program has three steps. The rest of this document keeps them separate.

Step Code What happens Side effects
Describe pipeline("events-topic") Builds a tree of instructions. None
Translate .foldMap(Interpreter) Changes each instruction into an IO action and chains them. None
Execute .unsafeRun() Runs the IO. Reads, writes, prints

What the demo pipeline does

  1. It reads events from a fake Kafka topic.
  2. It changes the events to upper case.
  3. It writes the events to a fake database.
  4. It reads all records from the database and prints the count.

The fake Kafka and the fake database fail 50 percent of the time. A Retryable instruction handles the failures.

Run the program

The program needs a JDK and sbt. The project sets sbt 1.10.7 in project/build.properties and Scala 3.3.5 in build.sbt. sbt downloads Scala.

sbt run

Main.scala runs the pipeline five times. Each run is random. Some runs show retries. Some runs fail after the last retry.

Concepts

Each concept has a definition, a diagram, and the code where it appears. Each concept uses the one before it.

ADT

ADT means algebraic data type. It does not mean "abstract data type". That is a different idea with the same initials.

An algebraic data type has two building blocks:

  • Sum: "one of these". DataOp is a sum. A value is a ReadFromKafka, or a TransformData, or a WriteToDatabase, and so on. In Scala, a sum is a sealed trait with case classes, or an enum.
  • Product: "all of these together". ReadFromKafka(topic: String) is a product with one field. Retryable(op, retries) is a product with two fields. In Scala, a product is a case class.

Two properties are important:

  1. The type is closed. sealed means that all cases are in this file. The compiler knows the full list. If a match in Interpreter misses a case, the compiler gives a warning.
  2. The type is only data. It has no methods and no behavior. ReadFromKafka("events") does not read. It is a note that says "read from events".

The A in DataOp[A] makes this a generalized algebraic data type, or GADT. Each case sets A to a different type. ReadFromKafka sets A to Either[String, List[String]]. LogMessage sets A to Unit. This tells the interpreter which type each instruction must return.

Functor

A functor F maps a category C to a category D. It maps each object a to F a. It maps each arrow f: a → b to F f: F a → F b. The structure does not change.

block-beta
  columns 3
  block:C:1
    columns 1
    hC["Category C"]
    a(("a"))
    space
    b(("b"))
  end
  space
  block:D:1
    columns 1
    hD["Category D"]
    Fa(("F a"))
    space
    Fb(("F b"))
  end

  a -- "f" --> b
  Fa -- "F f" --> Fb
  a -- "F" --> Fa
  b -- "F" --> Fb

  classDef header fill:none,stroke:none,color:#6F6A61
  classDef cat fill:#FBF7EE,stroke:#B8B2A6
  class hC,hD header
  class C,D cat
Loading

In Scala, the category is "Scala types" and the arrows are functions. A functor is a type constructor such as List, Option, or IO, together with map. A functor that maps a category to itself is an endofunctor. Every Scala type constructor is an endofunctor. It takes a Scala type and returns a Scala type.

Natural transformation

F and G are two functors from C to D. A natural transformation η gives one arrow η_a: F a → G a for every object a. The squares in the diagram commute: η_b ∘ F f = G f ∘ η_a. In words: map then transform gives the same result as transform then map.

block-beta
  columns 5
  block:C:1
    columns 1
    hC["Category C"]
    a(("a"))
    space
    space
    space
    space
    space
    b(("b"))
  end
  space
  block:D:3
    columns 3
    hD["Category D"]:3
    space:2 Ga(("G a"))
    space:3
    Fa(("F a")) space:2
    space:3
    Fb(("F b")) space:2
    space:3
    space:2 Gb(("G b"))
  end

  a -- "f" --> b
  Fa -- "F f" --> Fb
  Ga -- "G f" --> Gb
  a -- "F" --> Fa
  b -- "F" --> Fb
  a -- "G" --> Ga
  b -- "G" --> Gb
  Fa -- "η_a" --> Ga
  Fb -- "η_b" --> Gb

  classDef header fill:none,stroke:none,color:#6F6A61
  classDef cat fill:#FBF7EE,stroke:#B8B2A6
  class hC,hD header
  class C,D cat
Loading

Monad

pure puts a value into the monad. flatMap adds a step that returns a new monadic value.

flowchart TB
  a((a)) -- "pure(a)" --> MA(("M[A]"))
  MA -- "flatMap(f: A => M[B])" --> MB(("M[B]"))
  MB -- "flatMap(g: B => M[C])" --> MC(("M[C]"))
Loading

Free monad

An ADT becomes a free monad. The free monad gives us a DSL. foldMap sends the DSL to an interpreter. Two interpreters can run the same DSL.

flowchart TB
  ADT["ADT"]
  Free["Free monad"]
  DSL["DSL"]
  fold["foldMap(nt)"]
  I1["Interpreter 1"]
  I2["Interpreter 2"]

  ADT -- "Lift ADT into Free Monad" --> Free
  Free -- "DSL definition as computation using flatMap" --> DSL
  DSL -- "Transform using foldMap and apply Natural Transformation" --> fold
  fold -- "Interpreter 1 executes DSL" --> I1
  fold -- "Interpreter 2 executes DSL" --> I2
Loading

flatMap on Free does not run anything. It builds a FlatMap node. A Free program is a tree of data. The branches of the tree are functions. So the tree grows when values arrive.

The free monad does what a monad must do, pure and flatMap, and nothing more. It does not know what ReadFromKafka means. It does not know about Kafka, databases, or printing. It only knows "this step, then that step". The interpreter gives the meaning later.

foldMap

The DSL pipeline is on the left. The IO pipeline is on the right. Each instruction maps to one IO action. Each flatMap on the left maps to one flatMap on the right.

block-beta
  columns 3
  block:DSL:1
    columns 1
    hD["DSL Pipeline (Free Monad)"]
    r1["ReadFromKafka(topic)"]
    space
    l1["LogMessage('Read...')"]
    space
    t1["TransformData(records)"]
    space
    l2["LogMessage('Transformed...')"]
    space
    w1["WriteToDatabase(records)"]
  end
  space
  block:IO:1
    columns 1
    hI["Target Pipeline (IO)"]
    r2["Kafka.read(topic)"]
    space
    p1["println('Read...')"]
    space
    t2["Processor.transform(records)"]
    space
    p2["println('Transformed...')"]
    space
    w2["Database.write(records)"]
  end

  r1 -- "flatMap" --> l1
  l1 -- "flatMap" --> t1
  t1 -- "flatMap" --> l2
  l2 -- "flatMap" --> w1
  r2 -- "flatMap" --> p1
  p1 -- "flatMap" --> t2
  t2 -- "flatMap" --> p2
  p2 -- "flatMap" --> w2
  r1 -- "nt" --> r2
  l1 -- "nt" --> p1
  t1 -- "nt" --> t2
  l2 -- "nt" --> p2
  w1 -- "nt" --> w2

  classDef header fill:none,stroke:none,color:#6F6A61
  classDef cat fill:#FBF7EE,stroke:#B8B2A6
  classDef dsl fill:#FBEBDD,stroke:#D9731A,color:#1E1C19
  classDef io fill:#E3EEF7,stroke:#2F5D9A,color:#1E1C19
  class hD,hI header
  class DSL,IO cat
  class r1,l1,t1,l2,w1 dsl
  class r2,p1,t2,p2,w2 io
Loading

Orange boxes are DataOp instructions in the Free program. Blue boxes are the IO actions that the interpreter makes. Each nt arrow is one call to the natural transformation.

One run of the pipeline

One run has these steps:

  1. LogMessage with "Starting processing for topic".
  2. Retryable(ReadFromKafka(topic), 3). The interpreter calls retryOp. It makes a maximum of 3 attempts. Each failed attempt prints "Retrying due to error".
  3. If the read is successful: LogMessage with the records, TransformData (upper case), Retryable(WriteToDatabase(records), 2), then LogMessage with the result.
  4. If the read fails after 3 attempts: LogMessage with "Read from Kafka failed: Max retries reached".
  5. ReadAllFromDatabase. This prints the record count.
  6. LogMessage with a separator line.

Example output of one run. The event ids are random. The retries are random.

Starting processing for topic: events-topic
Reading from Kafka: events-topic
Retrying due to error: Failed to read from Kafka: events-topic
Reading from Kafka: events-topic
Read 2 records:
	event id:3f2a9c1e-... from events-topic
	event id:9c1d47b0-... from events-topic
Transforming records...
Transformed 2 records
Writing 2 records to the database.
Retrying due to error: Error writing to database
Writing 2 records to the database.
Write to DB successful
Processing complete
Database has 2 records

=============================

Resources

About

Free Monads for Declarative Data Pipelines

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Contributors

Languages