We want to write a pipeline like this:
for
records <- readFromKafka(topic)
cleaned <- transformData(records)
_ <- writeToDatabase(cleaned)
yield ()We also want two more things:
- 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.
- 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.
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 |
- It reads events from a fake Kafka topic.
- It changes the events to upper case.
- It writes the events to a fake database.
- 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.
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 runMain.scala runs the pipeline five times. Each run is random. Some runs show retries. Some runs fail after the last retry.
Each concept has a definition, a diagram, and the code where it appears. Each concept uses the one before it.
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".
DataOpis a sum. A value is aReadFromKafka, or aTransformData, or aWriteToDatabase, and so on. In Scala, a sum is asealed traitwith case classes, or anenum. - 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 acase class.
Two properties are important:
- The type is closed.
sealedmeans that all cases are in this file. The compiler knows the full list. If amatchinInterpretermisses a case, the compiler gives a warning. - 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.
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
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.
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
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]"))
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
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.
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
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 has these steps:
LogMessagewith "Starting processing for topic".Retryable(ReadFromKafka(topic), 3). The interpreter callsretryOp. It makes a maximum of 3 attempts. Each failed attempt prints "Retrying due to error".- If the read is successful:
LogMessagewith the records,TransformData(upper case),Retryable(WriteToDatabase(records), 2), thenLogMessagewith the result. - If the read fails after 3 attempts:
LogMessagewith "Read from Kafka failed: Max retries reached". ReadAllFromDatabase. This prints the record count.LogMessagewith 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
=============================
- The Interpreter Pattern Revisited
- Categories for the Working Hacker
- Functional Data Engineering: a modern paradigm for batch data processing
- Functional Data Engineering: A Set of Best Practices
- Free monads and event sourcing architecture
- Hexagonal Architecture and Free Monad: Two related design patterns?