-
Notifications
You must be signed in to change notification settings - Fork 46
/
WritePrimes.scala
69 lines (57 loc) · 2.2 KB
/
WritePrimes.scala
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
package sample.stream
import java.nio.file.Paths
import akka.NotUsed
import akka.actor.ActorSystem
import akka.stream.ActorMaterializer
import akka.stream.ClosedShape
import akka.stream.scaladsl._
import akka.util.ByteString
import scala.concurrent.forkjoin.ThreadLocalRandom
import scala.util.Failure
import scala.util.Success
object WritePrimes {
def main(args: Array[String]): Unit = {
implicit val system = ActorSystem("Sys")
import system.dispatcher
implicit val materializer = ActorMaterializer()
// generate random numbers
val maxRandomNumberSize = 1000000
val primeSource: Source[Int, NotUsed] =
Source.fromIterator(() => Iterator.continually(ThreadLocalRandom.current().nextInt(maxRandomNumberSize))).
// filter prime numbers
filter(rnd => isPrime(rnd)).
// and neighbor +2 is also prime
filter(prime => isPrime(prime + 2))
// write to file sink
val fileSink = FileIO.toPath(Paths.get("target/primes.txt"))
val slowSink = Flow[Int]
// act as if processing is really slow
.map(i => { Thread.sleep(1000); ByteString(i.toString) })
.toMat(fileSink)((_, bytesWritten) => bytesWritten)
// console output sink
val consoleSink = Sink.foreach[Int](println)
// send primes to both slow file sink and console sink using graph API
val graph = GraphDSL.create(slowSink, consoleSink)((slow, _) => slow) { implicit builder =>
(slow, console) =>
import GraphDSL.Implicits._
val broadcast = builder.add(Broadcast[Int](2)) // the splitter - like a Unix tee
primeSource ~> broadcast ~> slow // connect primes to splitter, and one side to file
broadcast ~> console // connect other side of splitter to console
ClosedShape
}
val materialized = RunnableGraph.fromGraph(graph).run()
// ensure the output file is closed and the system shutdown upon completion
materialized.onComplete {
case Success(_) =>
system.terminate()
case Failure(e) =>
println(s"Failure: ${e.getMessage}")
system.terminate()
}
}
def isPrime(n: Int): Boolean = {
if (n <= 1) false
else if (n == 2) true
else !(2 until n).exists(x => n % x == 0)
}
}