In [7]:
import $ivy.`com.typesafe.akka::akka-stream:2.6.4`
repl.pprinter() = repl.pprinter().copy(defaultHeight = 5 )

[32mimport [39m[36m$ivy.$                                     
[39m

In [8]:
import java.time._
import scala.concurrent._, duration._
import akka._
import akka.actor._
import akka.stream._
import akka.stream.scaladsl._

[32mimport [39m[36mjava.time._
[39m
[32mimport [39m[36mscala.concurrent._, duration._
[39m
[32mimport [39m[36makka._
[39m
[32mimport [39m[36makka.actor._
[39m
[32mimport [39m[36makka.stream._
[39m
[32mimport [39m[36makka.stream.scaladsl._[39m

# A streaming DSL

[Akka streams](https://doc.akka.io/api/akka/current/akka/stream/index.html) offers a DSL for programming reactive stream processors. These programs are made up from three basic components: sources, flows and sinks.

`Source`s represent data publishers. 

In [None]:
val source: Source[Int, NotUsed] = Source(List(1,2,3,4,5,6,7,8,9))

`Sink`s are data consumers. 

In [None]:
val sink: Sink[Any, Future[Done]]  = Sink.foreach(println)

We can connect a source and sink in order to obtain a so-called _runnable graph_, i.e. a streaming processor that can be actually run. In DSL terminology, the `RunnableGraph` is the _program_.

In [None]:
val graph1: RunnableGraph[NotUsed] = source.to(sink)

In order to run a graph we need a _materializer_, i.e. the interpreter of the streaming program. The standard materializer offered by akka-stream builts upon actors, so we need an actor system first.

In [9]:
implicit lazy val system: ActorSystem = ActorSystem("akka-stream-primer")

[36msystem[39m: [32mActorSystem[39m = [32m[lazy][39m

There is no need to explicitly instantiate the materializer, since it's already available implicitly:

In [10]:
implicitly[Materializer]

[36mres9[39m: [32mMaterializer[39m = [33mPhasedFusingActorMaterializer[39m(
  akka://akka-stream-primer,
  ActorMaterializerSettings(4,16,akka.actor.default-dispatcher,<function1>,StreamSubscriptionTimeoutSettings(CancelTermination,5000 milliseconds),false,1000,100...

We can now run our graph:

In [None]:
graph1.run

Between the sources and sinkes we can attach _flows_, intermediate processing steps: 

In [None]:
val graph2 = source.via(Flow[Int].map(_ + 1)).to(sink)

In [None]:
graph2.run

### Logging

We can log the activity of each operator in a graph to properly understand the contribution of each step in the transformation pipeline.

In [None]:
import akka.event.Logging

implicit class SourceOps[A, M](S: Source[A, M]){
    def logAll(l: String): Source[A, M] =
    S.log(l).withAttributes(Attributes.logLevels(
        onElement = Logging.WarningLevel,
        onFinish = Logging.WarningLevel,
        onFailure = Logging.DebugLevel))
}

In [None]:
Source(List(1,2,3)).logAll("source")
    .via(Flow[Int].map(_ + 1)).logAll("flow")
    .to(Sink.ignore)
    .run

# Materialized values

The output of a pipeline is called the *materialized value*.

In [61]:
val source = Source(List(1,2,3))

[36msource[39m: [32mSource[39m[[32mInt[39m, [32mNotUsed[39m] = Source(SourceShape(StatefulMapConcat.out(513068339)))

In [64]:
val ignoreS = Sink.ignore
val toListS = Sink.collection[Int,List[Int]]
val toPrintlnS = Sink.foreach(println)
val foldS = Sink.fold[String, Int]("")((acc: String, e: Int) => acc + e)

[36mignoreS[39m: [32mSink[39m[[32mAny[39m, [32mFuture[39m[[32mDone[39m]] = Sink(SinkShape(Ignore.in(1194690801)))
[36mtoListS[39m: [32mSink[39m[[32mInt[39m, [32mFuture[39m[[32mList[39m[[32mInt[39m]]] = Sink(SinkShape(seq.in(1807803460)))
[36mtoPrintlnS[39m: [32mSink[39m[[32mAny[39m, [32mFuture[39m[[32mDone[39m]] = Sink(SinkShape(Map.in(39437523)))
[36mfoldS[39m: [32mSink[39m[[32mInt[39m, [32mFuture[39m[[32mString[39m]] = Sink(SinkShape(Fold.in(1608103017)))

In [65]:
source.to(ignoreS)
source.to(toListS)
source.to(foldS)

[36mres64_0[39m: [32mRunnableGraph[39m[[32mNotUsed[39m] = [33mRunnableGraph[39m(
  [33mLinearTraversalBuilder[39m(
    None,
    None,
...
[36mres64_1[39m: [32mRunnableGraph[39m[[32mNotUsed[39m] = [33mRunnableGraph[39m(
  [33mLinearTraversalBuilder[39m(
    None,
    None,
...
[36mres64_2[39m: [32mRunnableGraph[39m[[32mNotUsed[39m] = [33mRunnableGraph[39m(
  [33mLinearTraversalBuilder[39m(
    None,
    None,
...

In [66]:
s.toMat(ignoreS)(Keep.right)
s.toMat(toListS)(Keep.right)
s.toMat(foldS)(Keep.right)

[36mres65_0[39m: [32mRunnableGraph[39m[[32mFuture[39m[[32mDone[39m]] = [33mRunnableGraph[39m(
  [33mLinearTraversalBuilder[39m(
    None,
    None,
...
[36mres65_1[39m: [32mRunnableGraph[39m[[32mFuture[39m[[32mList[39m[[32mInt[39m]]] = [33mRunnableGraph[39m(
  [33mLinearTraversalBuilder[39m(
    None,
    None,
...
[36mres65_2[39m: [32mRunnableGraph[39m[[32mFuture[39m[[32mString[39m]] = [33mRunnableGraph[39m(
  [33mLinearTraversalBuilder[39m(
    None,
    None,
...

Common shortcuts:

# Async boundaries 

### Akka streams vs. iterators

Which is the difference between the previous akka stream program and the following `Iterator` program?

In [None]:
List(1,2,3,4).iterator.map(_ + 1).foreach(println)

In both cases, we obtain a streaming processor. However, in akka streams, intermediate processing steps and actions performed over the resulting data, are first-class entities: flows and sinks. They can be defined independently, reused and combined as we wish. There is no such notion in the iterator realm. Moreover, akka streams are compiled into actors, and the graph has the potential to run asyncronously and concurrently, with back-pressure niceties. Iterator programs are run sequentially and syncronously.  

### Exploiting parallelism 

In [None]:
Source(1 to 3)
    .map { i =>
        println(s"A: $i"); i
    }
    .async
    .map { i =>
        println(s"B: $i"); i
    }
    .async
    .map { i =>
        println(s"C: $i"); i
    }
    .to(Sink.ignore)
    .run

# Fan-in, fan-out & additional operators

In [11]:
val source1: Source[Int, NotUsed] = Source(List(1,2,3,4))
val source2: Source[String, NotUsed] = Source(List("hola", "adios"))

[36msource1[39m: [32mSource[39m[[32mInt[39m, [32mNotUsed[39m] = Source(SourceShape(StatefulMapConcat.out(1243834587)))
[36msource2[39m: [32mSource[39m[[32mString[39m, [32mNotUsed[39m] = Source(SourceShape(StatefulMapConcat.out(529404463)))

In [13]:
source1.map(i => s"num: $i")
    .merge(source2)
    .runWith(Sink.foreach(println))

num: 1
hola
num: 2
adios
num: 3
num: 4


In [14]:
source1.zip(source2)
    .runWith(Sink.foreach(println))

(1,hola)
(2,adios)


In [15]:
source1.map(_.toString)
    .concat(source2)
    .runWith(Sink.foreach(println))

1
2
3
4
hola
adios


See https://doc.akka.io/docs/akka/current/stream/stream-substream.html, for an explanation of substreams.

# File IO

In [29]:
import java.nio.file._, akka.util._

[32mimport [39m[36mjava.nio.file._, akka.util._[39m

In [30]:
val file = Paths.get("Intro.ipynb")

[36mfile[39m: [32mPath[39m = Intro.ipynb

In [47]:
val lines: Source[String, Future[IOResult]] = FileIO.fromPath(file)
    .via(Framing.delimiter(ByteString("\n"), maximumFrameLength = 1500, allowTruncation = true))
    .throttle(1, 10.millisecond)
    .map(_.utf8String)
//    .runWith(Sink.takeLast(3))
//    .runWith(Sink.foreach(println))

[36mlines[39m: [32mSource[39m[[32mString[39m, [32mFuture[39m[[32mIOResult[39m]] = Source(SourceShape(Map.out(1086145708)))

In [48]:
lines.runWith(Sink.foreach(println))

{
 "cells": [
  {
   "cell_type": "code",
   "execution_count": 7,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "text/plain": [


In [52]:
lines.map(_.length).map(_.toString+"\n").map(ByteString(_)).runWith(FileIO.toPath(Paths.get("linecounts.txt")))

In [50]:
Source(List("a","b","c"))
    .map(ByteString(_))
    .runWith(FileIO.toPath(Paths.get("linecounts.txt")))