This repository was archived by the owner on Sep 23, 2019. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 24
Expand file tree
/
Copy pathWorkPullingPattern.scala
More file actions
80 lines (69 loc) · 2.24 KB
/
Copy pathWorkPullingPattern.scala
File metadata and controls
80 lines (69 loc) · 2.24 KB
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
70
71
72
73
74
75
76
77
78
79
80
package akkapatterns
import akka.actor.Actor
import scala.collection.mutable
import akka.actor.ActorRef
import WorkPullingPattern._
import scala.collection.IterableLike
import scala.reflect.ClassTag
import org.slf4j.LoggerFactory
import akka.actor.Terminated
object WorkPullingPattern {
sealed trait Message
trait Epic[T] extends Iterable[T] //used by master to create work (in a streaming way)
case object GimmeWork extends Message
case object CurrentlyBusy extends Message
case object WorkAvailable extends Message
case class RegisterWorker(worker: ActorRef) extends Message
case class Work[T](work: T) extends Message
}
class Master[T] extends Actor {
val log = LoggerFactory.getLogger(getClass)
val workers = mutable.Set.empty[ActorRef]
var currentEpic: Option[Epic[T]] = None
def receive = {
case epic: Epic[T] ⇒
if (currentEpic.isDefined)
sender ! CurrentlyBusy
else if (workers.isEmpty)
log.error("Got work but there are no workers registered.")
else {
currentEpic = Some(epic)
workers foreach { _ ! WorkAvailable }
}
case RegisterWorker(worker) ⇒
log.info(s"worker $worker registered")
context.watch(worker)
workers += worker
case Terminated(worker) ⇒
log.info(s"worker $worker died - taking off the set of workers")
workers.remove(worker)
case GimmeWork ⇒ currentEpic match {
case None ⇒
log.info("workers asked for work but we've no more work to do")
case Some(epic) ⇒
val iter = epic.iterator
if (iter.hasNext)
sender ! Work(iter.next)
else {
log.info(s"done with current epic $epic")
currentEpic = None
}
}
}
}
abstract class Worker[T](val master: ActorRef) extends Actor {
override def preStart {
master ! RegisterWorker(self)
master ! GimmeWork
}
def receive = {
case WorkAvailable ⇒
master ! GimmeWork
case Work(work: T) ⇒
// haven't found a nice way to get rid of that warning
// looks like we can't suppress the erasure warning: http://stackoverflow.com/questions/3506370/is-there-an-equivalent-to-suppresswarnings-in-scala
doWork(work)
master ! GimmeWork
}
def doWork(work: T)
}