This is an old revision of the document!
아래 과정은 console에서 수행할 경우 필요한 작업. Intellij로 sbt 프로젝트 생성하면 더 쉽게 만들 수 있다.
set name := "hello-akka" set version := "0.1" set scalaVersion := "2.12.6" session save exit
// https://mvnrepository.com/artifact/com.typesafe.akka/akka-actor libraryDependencies += "com.typesafe.akka" %% "akka-actor" % "2.5.14"
package com.practice.ex01
import akka.actor.ActorSystem
object HelloAkkaActorSystem extends App{
val actorSystem = ActorSystem("HelloAkka")
println(actorSystem)
// akka://HelloAkka
}
sbt "runMain com.practice.ex01.HelloAkkaActorSystem"
scala에서 ask는 “?”, tell은 “!” method를 사용한다.
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import akka.pattern.ask
import akka.util.Timeout
import scala.concurrent.Await
import scala.concurrent.duration._
import akka.actor.Actor
object HelloAkkaActorSystem extends App {
val actorSystem = ActorSystem("HelloAkka")
val actor = actorSystem.actorOf(Props(classOf[SumActor]), "sumActor")
// ask를 사용하려면 timeout을 반드시 명시해야 함.
implicit val timeout = Timeout(10 seconds)
val future = actor ? (1 to 10).toArray
val result = Await.result(future, 10 seconds)
println(s"result : $result")
}
class SumActor extends Actor {
override def receive: Receive = {
case arr: Array[Int] => sender ! arr.sum
case _ => sender ! "Error!"
}
}
ask는 기본적으로 blocking. 그러므로 비슷하게 tell을 이용하여 서로 주고 받는 callback 형식으로 구현하므로 더 경제적.
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import com.practice.ex01.Message.Start
object HelloAkkaActorSystem extends App {
val actorSystem = ActorSystem("HelloAkka")
val queryActor = actorSystem.actorOf(Props(classOf[QueryActor]), "queryActor")
val sumActor = actorSystem.actorOf(Props(classOf[SumActor]), "sumActor")
queryActor ! Start(sumActor, (1 to 10).toArray)
}
package com.practice.ex01
import akka.actor.Actor
import com.practice.ex01.Message._
class QueryActor extends Actor{
override def receive: Receive = {
case Start(actorRef, arr) =>
println("Start!")
actorRef ! Progress(arr)
case Complete(result) =>
println(s"Complete! $result")
case _ => println("QueryActor - WTHIGO?")
}
}
package com.practice.ex01
import akka.actor.Actor
import com.practice.ex01.Message._
class SumActor extends Actor {
override def receive: Receive = {
case Progress(arr) =>
println("Progress!")
sender ! Complete(arr.sum)
case _ => sender ! "SumActor - WTHIGO?"
}
}
Actor가 메시지를 가져오는 방법을 제어할 때 사용.
package com.practice.ex01
import akka.actor.{ActorRef, ActorSystem}
import akka.dispatch.{MailboxType, MessageQueue, ProducesMessageQueue}
import com.typesafe.config.Config
class MyCustomMailbox extends MailboxType with ProducesMessageQueue[MyMessageQueue] {
def this(setting: ActorSystem.Settings, config: Config) = this()
override def create(owner: Option[ActorRef], system: Option[ActorSystem]): MessageQueue = {
println("MyCustomMailbox is created!")
new MyMessageQueue()
}
}
package com.practice.ex01
import java.util.concurrent.ConcurrentLinkedQueue
import akka.actor.ActorRef
import akka.dispatch.{Envelope, MessageQueue}
class MyMessageQueue extends MessageQueue{
private final val queue = new ConcurrentLinkedQueue[Envelope]()
override def enqueue(receiver: ActorRef, handle: Envelope): Unit = {
println("enqueue!")
if(handle.sender.path.name == "MyActor") {
handle.sender ! "I know you!"
queue.offer(handle)
}
else
handle.sender ! "Who are you?"
}
override def dequeue(): Envelope = queue.poll
override def numberOfMessages: Int = queue.size
override def hasMessages: Boolean = !queue.isEmpty
override def cleanUp(owner: ActorRef, deadLetters: MessageQueue): Unit = {
while (hasMessages) {
deadLetters.enqueue(owner, dequeue())
}
}
}
custom-dispatcher {
mailbox-requirement = "com.practice.ex01.MyMessageQueue"
}
akka.actor.mailbox.requirements {
"com.practice.ex01.MyMessageQueue" = custom-dispatcher-mailbox
}
custom-dispatcher-mailbox {
mailbox-type = "com.practice.ex01.MyCustomMailbox"
}
package com.practice.ex01
import akka.actor.Actor
class MySpecialActor extends Actor {
override def receive: Receive = {
case msg: String => println(s"MySpecialActor's msg : $msg")
}
}
package com.practice.ex01
import akka.actor.{Actor, ActorRef}
class MyActor extends Actor {
override def receive: Receive = {
case (msg: String, actorRef: ActorRef) => actorRef ! msg
case msg => println(s"MyActor's msg? $msg")
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
object Client extends App {
val actorSystem = ActorSystem("HelloAkka")
// application.conf에 기술된 custom-dispatcher을 읽고
// 사용자 정의 Mailbox(MyCustomMailBox)가 먼저 생성되고,
// 그 안에 create() method로 message queue가 생성된다.
val mySpecialActor = actorSystem.actorOf(Props[MySpecialActor]
.withDispatcher("custom-dispatcher"))
val unknownActor = actorSystem.actorOf(Props[MyActor], "abracadabra")
val knownActor = actorSystem.actorOf(Props[MyActor], "MyActor")
unknownActor ! ("hello", mySpecialActor)
knownActor ! ("hello", mySpecialActor)
// 비동기적이므로 아래 결과는 바뀔 수 있음.
/*
MyCustomMailbox is created!
enqueue!
enqueue!
MyActor's msg? Who are you?
MyActor's msg? I know you!
MySpecialActor's msg : hello
*/
}
우선순위에 따라 특정 메시지를 먼저 처리하고 싶을 때 사용. 설정은 위와 비슷하므로 상세한 설명은 생략.
package com.practice.ex01
import akka.actor.ActorSystem
import akka.dispatch.{PriorityGenerator, UnboundedPriorityMailbox}
import com.typesafe.config.Config
class MyPriorityMailbox(settings: ActorSystem.Settings, config: Config)
extends UnboundedPriorityMailbox (
PriorityGenerator {
// int 값이 작을수록 우선순위가 높다.
case x: Int => 1
case x: String => 0
case _ => 3
}
)
my-priority-mailbox {
mailbox-type = "com.practice.ex01.MyPriorityMailbox"
}
package com.practice.ex01
import akka.actor.Actor
class MyPriorityActor extends Actor{
override def receive: Receive = {
case x: Int => println(x)
case x: String => println(x)
case x => println(x)
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
object Client extends App {
val actorSystem = ActorSystem("HelloAkka")
val myPriorityActor = actorSystem.actorOf(Props[MyPriorityActor]
.withDispatcher("my-priority-mailbox"))
Array('LowestPriority, 123, "Top Priority!", 1.23).foreach(myPriorityActor ! _)
/* 결과적으로 String, int 순으로 우선순위가 높음. */
/*
Top Priority!
123
1.23
'LowestPriority
*/
}
package com.practice.ex01 import akka.dispatch.ControlMessage case object MyControlMessage extends ControlMessage
control-aware-mailbox {
mailbox-type = "akka.dispatch.UnboundedControlAwareMailbox"
}
package com.practice.ex01
import akka.actor.Actor
class Logger extends Actor{
override def receive: Receive = {
case MyControlMessage => println("Control message first!")
case x => println(x)
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
object Client extends App {
val actorSystem = ActorSystem("HelloAkka")
val loggerActor = actorSystem.actorOf(Props[Logger]
.withDispatcher("control-aware-mailbox"))
loggerActor ! "Hello,"
loggerActor ! "world!"
loggerActor ! MyControlMessage // 우선적으로 처리
/*
Control message first!
Hello,
world!
*/
}
become/unbecome으로 구현.
package com.practice.ex01
object StateMessage {
case class IntState()
case class StringState()
}
package com.practice.ex01
import akka.actor.Actor
import com.practice.ex01.StateMessage.{IntState, StringState}
class StateChangingActor extends Actor {
override def receive: Receive = {
case StringState => context.become(isStateString)
case IntState => context.become(isStateInt)
}
def isStateString: Receive = {
case msg : String => println(s"$msg")
case IntState => context.become(isStateInt)
}
def isStateInt: Receive = {
case msg : Int => println(s"$msg")
case StringState => context.become(isStateString)
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import com.practice.ex01.StateMessage.{IntState, StringState}
object Client extends App {
val actorSystem = ActorSystem("HelloAkka")
val stateChangingActor = actorSystem.actorOf(Props[StateChangingActor])
stateChangingActor ! StringState
stateChangingActor ! "Hello, world!"
stateChangingActor ! IntState
stateChangingActor ! 123
stateChangingActor ! "abracadabra" // 현재 state가 IntState이기 때문에 무시됨.
/*
Hello, world!
123
*/
}
https://stackoverflow.com/questions/13847963/akka-kill-vs-stop-vs-poison-pill
package com.practice.ex01
import akka.actor.Actor
class ShutdownActor extends Actor {
override def receive: Receive = {
case msg: String => println(s"$msg")
case Stop => context.stop(self)
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Kill, PoisonPill, Props}
object Client extends App {
val actorSystem = ActorSystem("HelloAkka")
val shutdownActor1 = actorSystem.actorOf(Props[ShutdownActor])
shutdownActor1 ! "hello1"
shutdownActor1 ! PoisonPill
shutdownActor1 ! "Are you breathing?"
val shutdownActor2 = actorSystem.actorOf(Props[ShutdownActor])
shutdownActor2 ! "hello2"
shutdownActor2 ! Kill
shutdownActor2 ! "Are you breathing?"
val shutdownActor3 = actorSystem.actorOf(Props[ShutdownActor])
shutdownActor3 ! "hello3"
shutdownActor3 ! Stop
shutdownActor3 ! "Are you breathing?"
/*
hello2
hello1
hello3
scala> [INFO] [08/06/2018 00:25:30.860] [HelloAkka-akka.actor.default-dispatcher-2] [akka://HelloAkka/user/$c] Message [java.lang.String] without sender to Actor[akka://HelloAkka/user/$c#303943664] was not delivered. [1] dead letters encountered. If this is not an expected behavior, then [Actor[akka://HelloAkka/user/$c#303943664]] may have terminated unexpectedly, This logging can be turned off or adjusted with configuration settings 'akka.log-dead-letters' and 'akka.log-dead-letters-during-shutdown'.
[INFO] [08/06/2018 00:25:30.861] [HelloAkka-akka.actor.default-dispatcher-2] [akka://HelloAkka/user/$a] Message [java.lang.String] without sender to Actor[akka://HelloAkka/user/$a#-1903145472] was not delivered. [2] dead letters encountered. If this is not an expected behavior, then [Actor[akka://HelloAkka/user/$a#-1903145472]] may have terminated unexpectedly, This logging can be turned off or adjusted with configuration settings 'akka.log-dead-letters' and 'akka.log-dead-letters-during-shutdown'.
[ERROR] [08/06/2018 00:25:30.862] [HelloAkka-akka.actor.default-dispatcher-5] [akka://HelloAkka/user/$b] Kill (akka.actor.ActorKilledException: Kill)
[INFO] [08/06/2018 00:25:30.863] [HelloAkka-akka.actor.default-dispatcher-2] [akka://HelloAkka/user/$b] Message [java.lang.String] without sender to Actor[akka://HelloAkka/user/$b#14529698] was not delivered. [3] dead letters encountered. If this is not an expected behavior, then [Actor[akka://HelloAkka/user/$b#14529698]] may have terminated unexpectedly, This logging can be turned off or adjusted with configuration settings 'akka.log-dead-letters' and 'akka.log-dead-letters-during-shutdown'.
*/
}
장애 발생시 정지되는 것이 아닌 처리량이 감소한 가동으로 항상 반응성을 유지하며, 완전 가동중과 비교하여 정책적으로 더 혹은 덜 가동되는 시스템.
특정 서비스가 살아있는 지 지속적으로 확인할 때 필요.
// 감시 context.watch(childActor: ActorRef) // 감시 해제 context.unwatch(childActor: ActorRef)
아래와 같은 Tree 구조로 운영한 뒤 Strategy, watch(감시, 모니터링)를 이용한다.
Master - Slave. 자식 액터의 장애를 부모가 처리하는 구조.
package com.practice.ex01
object Message {
case object CreateChild
case class Greet(msg: String)
}
package com.practice.ex01
import akka.actor.Actor
import com.practice.ex01.Message.Greet
class ChildActor extends Actor {
override def receive: Receive = {
case Greet(msg) => println(s"parent : ${self.path.parent} // me : ${self.path} // msg : ${msg}")
}
}
package com.practice.ex01
import akka.actor.{Actor, Props}
import com.practice.ex01.Message.{CreateChild, Greet}
class ParentActor extends Actor{
override def receive: Receive = {
case CreateChild =>
val child = context.actorOf(Props[ChildActor], "child")
child ! Greet("Hello, child!")
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import com.practice.ex01.Message.CreateChild
object Client extends App {
val actorSystem = ActorSystem("Supervision")
val parent = actorSystem.actorOf(Props[ParentActor], "parent")
parent ! CreateChild
// parent : akka://Supervision/user/parent // me : akka://Supervision/user/parent/child // msg : Hello, child!
}
https://doc.akka.io/docs/akka/2.5/routing.html
언제 사용하는가?
메시지 수가 가장 적은 액터에 메시지 전달. 즉, 가장 덜 바쁜 액터에게 메시지 전달.
class RoutingActor extends Actor {
override def receive: Receive = {
case msg: String => sender ! s"I am ${self.path.name}, I received $msg"
case _ => println(s"I don't understand the message")
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import akka.routing.SmallestMailboxPool
object Client extends App {
val actorSystem = ActorSystem("Supervision");
val router = actorSystem.actorOf(SmallestMailboxPool(5).props(Props[SmallestMailboxActor]))
1 to 5 foreach(i => router ! s"Hello $i")
/*
I am $a
I am $a
I am $d
I am $b
I am $c
*/
}
액터 사이의 작업을 재분배.
https://groups.google.com/forum/#!msg/akka-user/pymtrCZBduk/rwXKIdGUGvEJ https://doc.akka.io/docs/akka/2.5/routing.html#balancing-pool
각 액터에게 하나씩 메시지 전달.
말그대로 무작위.
하나의 같은 메시지를 모든 액터에게 전달. 모든 액터에게 같은 작업을 하도록 일반적인 명령을 보내고 싶을 때 사용.
https://doc.akka.io/docs/akka/2.5/routing.html#specially-handled-messages
액터를 관리하는 데 사용.
같은 라우터를 공유하는 액터들 중에게 같은 메시지를 모두 보내고(Broadcast), 가장 먼저 작업을 마친 액터를 기다려 응답을 송신(ask) 후 나머지 액터 응답을 버림. 여러 서버 중 가장 빠르게 반응하는 서버에 작업을 보내는 상황에서 사용.
package com.practice.ex01
import akka.actor.Actor
class RoutingActor extends Actor {
override def receive: Receive = {
case msg: String => sender ! s"I am ${self.path.name}, I received $msg"
case _ => println(s"I don't understand the message")
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import akka.pattern.ask
import akka.routing.ScatterGatherFirstCompletedPool
import akka.util.Timeout
import scala.concurrent.duration._
import scala.concurrent.{Await, Future}
object Client extends App {
val duration = 10 seconds
implicit val timeout = Timeout(duration)
val actorSystem = ActorSystem("Hello-Akka");
val router = actorSystem
.actorOf(ScatterGatherFirstCompletedPool(5, within = duration)
.props(Props[RoutingActor]))
val future: Future[Any] = router ? "hello"
val response: String = Await.result(future.mapTo[String], duration)
println(response)
// I am $a, I received hello
}
무작위로 고른 액터에게 메시지를 보내고, interval에서 지정한 약간의 지연 후 남은 액터들 중에 무작위로 고른 액터에게 송신하는 일을 계속한다. 첫 번째 응답이 수신되기를 기다리고, 이를 본래 송신자에게 전달, 나머지 응답은 버림.
package com.practice.ex01
import akka.actor.Actor
class RoutingActor extends Actor {
override def receive: Receive = {
case msg: String => sender ! s"I am ${self.path.name}, I received $msg"
case _ => println(s"I don't understand the message")
}
}
package com.practice.ex01
import akka.actor.{ActorSystem, Props}
import akka.pattern.ask
import akka.routing.TailChoppingPool
import akka.util.Timeout
import scala.concurrent.duration._
import scala.concurrent.{Await, Future}
object Client extends App {
val duration = 10 seconds
implicit val timeout = Timeout(duration)
val actorSystem = ActorSystem("Hello-Akka");
val router = actorSystem
.actorOf(TailChoppingPool(5, within = duration, interval = 20 millis)
.props(Props[RoutingActor]))
val future: Future[Any] = router ? "hello!"
val response: String = Await.result(future.mapTo[String], duration)
println(response)
}
말그대로 Hashing. Consistent Hashing 사용. 같은 키를 가지는 메시지를 항상 같은 액터에게 전달.
package com.practice.ex01
import akka.actor.Actor
import com.practice.ex01.Message._
class Cache extends Actor {
var cache = Map.empty[String, String]
override def receive: Receive = {
case Create(key, value) =>
println(s"${self.path.name} : Create $key/$value")
cache += (key -> value)
case Get(key) =>
println(s"${self.path.name} : Get $key")
context.sender ! cache.get(key)
case Remove(key) =>
println(s"${self.path.name} : Remove $key")
cache -= key
}
}
package com.practice.ex01
import akka.actor.{ActorRef, ActorSystem, Props}
import akka.pattern.ask
import akka.routing.ConsistentHashingPool
import akka.routing.ConsistentHashingRouter.{ConsistentHashMapping, ConsistentHashableEnvelope}
import akka.util.Timeout
import com.practice.ex01.Message._
import scala.concurrent.Await
import scala.concurrent.duration._
object Client extends App {
def hashMapping: ConsistentHashMapping = {
case Remove(key) ⇒ key
}
val actorSystem = ActorSystem("Hello-Akka")
val cache: ActorRef =
actorSystem.actorOf(
ConsistentHashingPool(10, hashMapping = hashMapping).
props(Props[Cache]), name = "cache")
// 생성
cache ! ConsistentHashableEnvelope(message = Create("hello", "HELLO"), hashKey = "hello")
cache ! ConsistentHashableEnvelope(message = Create("hi", "HI"), hashKey = "hi")
// Get 응답 받기, 삭제등 시험
val duration = 10 seconds
implicit val timeout = Timeout(duration)
val future1 = cache ? Get("hello")
val response1 = Await.result(future1.mapTo[Any], duration)
println(response1)
val future2 = cache ? Get("hi")
val response2 = Await.result(future2.mapTo[Any], duration)
println(response2)
cache ! Remove("hi")
val future3 = cache ? Get("hi")
val response3 = Await.result(future3.mapTo[Any], duration)
println(response3)
}
$d : Create hello/HELLO $e : Create hi/HI $d : Get hello Some(HELLO) $e : Get hi Some(HI) $e : Remove hi $e : Get hi None
https://doc.akka.io/docs/akka/2.5/routing.html#specially-handled-messages https://doc.akka.io/docs/akka/2.1.2/scala/routing.html
package com.practice.ex01
import akka.actor.{Actor, ActorSystem, Props}
import akka.pattern.ask
import akka.routing.{DefaultResizer, ScatterGatherFirstCompletedPool}
import akka.util.Timeout
import scala.annotation.tailrec
import scala.concurrent.duration._
import scala.concurrent.{Await, Future}
// Message
case class FibonacciNumber(nbr: Int)
// Actor
class FibonacciActor extends Actor {
def receive = {
case FibonacciNumber(nbr) => sender ! fibonacci(nbr)
case _ => new IllegalArgumentException
}
private def fibonacci(n: Int): Int = {
@tailrec
def fib(n: Int, b: Int, a: Int): Int = n match {
case 0 => a
case _ => fib(n - 1, a + b, b)
}
fib(n, 1, 0)
}
}
// Client
object Client extends App {
val duration = 10 seconds
implicit val timeout = Timeout(duration)
val akkaSystem = ActorSystem("Hello-Akka")
val resizer = DefaultResizer(lowerBound = 4, upperBound = 15)
val router = akkaSystem.actorOf(
ScatterGatherFirstCompletedPool(5, within = duration, resizer = Some(resizer))
.props(Props[FibonacciActor]))
val future: Future[Any] = router ? FibonacciNumber(20)
val response: Int = Await.result(future.mapTo[Int], duration)
println(response)
//6765
}