maven으로 scala AKKA echo 서버 구현하기
자바, 스칼라의 오픈소스 배포용으로 메이븐이 많이 사용되고 있다. 메이븐을 이용한 스칼라 프로젝트 사용법을 정리하고자 이 포스트를 작성한다.
소스 프로젝트는 깃허브로 배포해두었다. https://github.com/jamsya/akka-echo
- 메이븐 프로젝트 생성 삽질
메이븐에서 자동으로 생성하는 프로젝트로는 컴파일이 되지 않았다.
$ mvn archetype:generate
스칼라 버전 바꿔보고, 지웠다 깔고, jdk 바꿔보고 다 해봤는데 소용 없었다.
헌데 git에서 예제로 받은 코드는 메이븐으로 컴파일이 잘 된다 -_-;; 지금 메이븐에서의 스칼라 프로젝트 생성에는 문제가 있어보인다. 아니면 내가 모르는게 또 있거나..
그래서 깃헙에 있는 maven 샘플을 받아 사용하기로 했다.
- 표준 maven 샘플 받아오기
https://github.com/davidB/scala-archetype-simple
여기서 maven-simple 프로젝트를 사용했다.
/usr/local/src/scala-module-dependency-sample/maven-sample 이 디렉토리를 앞으로 계속 복붙해 쓰게 될 것이다.
복붙한 maven-sample 프로젝트는 다음과 같이 구성되어 있다. java를 추가하고 싶다면 main 디렉토리 밑에 java 디렉토리를 만들어 넣으면 된다.
* pom.xml * src * main * scala * com/my-package/… *.scala * test * scala * com/my-package/… *.scala * target … src/main/scala 디렉토리에 필요한 scala 소스코드를 넣는다. 에코서버 소스 코드는 포스트 맨 밑에 첨부해둔다.
- 컴파일 & 설치
$mvn compile
$mvn install
- 실행
메이븐 프로젝트에서의 스칼라 프로그램 실행
$ mvn scala:run -DmainClass=$메인클래스명 ## ex) mvn scala:run -DmainClass=package.EchoServer
- 프로필 바꾸기
메이븐에서는 각기 다른 시스템에 맞게 프로필을 사용하는 것이 가능하다. 특정 프로필로 컴파일하고 싶을 때 다음과 같이 옵션을 추가한다.
mvn compile -Pscala-2.11
- 에코서버 소스코드
EchoServer.scala
import java.net.InetSocketAddress
import akka.actor.{Actor, ActorDSL, ActorLogging, ActorRef, ActorSystem, Props, SupervisorStrategy} import akka.io.{IO, Tcp} import akka.util.ByteString import com.typesafe.config.ConfigFactory
import scala.concurrent.duration.DurationInt
object EchoServer extends App {
val config = ConfigFactory.parseString(“akka.loglevel = DEBUG”) implicit val system = ActorSystem(“EchoServer”, config)
// make sure to stop the system so that the application stops try run() finally system.terminate()
def run(): Unit = { import ActorDSL._
// create two EchoManager and stop the application once one dies val watcher = inbox() watcher.watch(system.actorOf(Props(classOf[EchoManager], classOf[EchoHandler]), “echo”)) watcher.watch(system.actorOf(Props(classOf[EchoManager], classOf[SimpleEchoHandler]), “simple”)) watcher.receive(10.minutes) }
}
class EchoManager(handlerClass: Class[_]) extends Actor with ActorLogging {
import Tcp._ import context.system
// there is not recovery for broken connections override val supervisorStrategy = SupervisorStrategy.stoppingStrategy
// bind to the listen port; the port will automatically be closed once this actor dies override def preStart(): Unit = { IO(Tcp) ! Bind(self, new InetSocketAddress(“localhost”, 0)) }
// do not restart override def postRestart(thr: Throwable): Unit = context stop self
def receive = { case Bound(localAddress) ⇒ log.info(“listening on port {}”, localAddress.getPort)
case CommandFailed(Bind(_, local, _, _, _)) ⇒ log.warning(s”cannot bind to [$local]”) context stop self
//#echo-manager case Connected(remote, local) ⇒ // log.info(“received connection from {}”, remote) val handler = context.actorOf(Props(handlerClass, sender(), remote)) sender() ! Register(handler, keepOpenOnPeerClosed = true) //#echo-manager }
}
//#echo-handler object EchoHandler { final case class Ack(offset: Int) extends Tcp.Event
def props(connection: ActorRef, remote: InetSocketAddress): Props = Props(classOf[EchoHandler], connection, remote) }
class EchoHandler(connection: ActorRef, remote: InetSocketAddress) extends Actor with ActorLogging {
import EchoHandler._ import Tcp._
// sign death pact: this actor terminates when connection breaks context watch connection
// start out in optimistic write-through mode def receive = writing
//#writing def writing: Receive = { case Received(data) ⇒ connection ! Write(data, Ack(currentOffset)) buffer(data) context become closing
// context stop self //여기서 닫으면 되나?? // context become closing;
case Ack(ack) ⇒ acknowledge(ack)
case CommandFailed(Write(_, Ack(ack))) ⇒ connection ! ResumeWriting context become buffering(ack)
case PeerClosed ⇒ if (storage.isEmpty) context stop self else context become closing } //#writing
//#buffering def buffering(nack: Int): Receive = { var toAck = 10 var peerClosed = false
{ case Received(data) ⇒ buffer(data) case WritingResumed ⇒ writeFirst() case PeerClosed ⇒ peerClosed = true case Ack(ack) if ack < nack ⇒ acknowledge(ack) case Ack(ack) ⇒ acknowledge(ack) if (storage.nonEmpty) { if (toAck > 0) { // stay in ACK-based mode for a while writeFirst() toAck -= 1 } else { // then return to NACK-based again writeAll() context become (if (peerClosed) closing else writing) } } else if (peerClosed) context stop self else context become writing } } //#buffering
//#closing def closing: Receive = { case CommandFailed(_: Write) ⇒ connection ! ResumeWriting context.become({
case WritingResumed ⇒ writeAll() context.unbecome()
case ack: Int ⇒ acknowledge(ack)
}, discardOld = false)
case Ack(ack) ⇒ acknowledge(ack) if (storage.isEmpty) context stop self } //#closing
override def postStop(): Unit = { log.info(s”transferred $transferred bytes from/to [$remote]”) }
//#storage-omitted private var storageOffset = 0 private var storage = Vector.empty[ByteString] private var stored = 0L private var transferred = 0L
val maxStored = 100000000L val highWatermark = maxStored * 5 / 10 val lowWatermark = maxStored * 3 / 10 private var suspended = false
private def currentOffset = storageOffset + storage.size
//#helpers private def buffer(data: ByteString): Unit = { storage :+= data stored += data.size
if (stored > maxStored) { log.warning(s”drop connection to [$remote] (buffer overrun)”) context stop self
} else if (stored > highWatermark) { log.debug(s”suspending reading at $currentOffset”) connection ! SuspendReading suspended = true } }
private def acknowledge(ack: Int): Unit = { require(ack == storageOffset, s”received ack $ack at $storageOffset”) require(storage.nonEmpty, s”storage was empty at ack $ack”)
val size = storage(0).size stored -= size transferred += size
storageOffset += 1 storage = storage drop 1
if (suspended && stored < lowWatermark) { log.debug(“resuming reading”) connection ! ResumeReading suspended = false } } //#helpers
private def writeFirst(): Unit = { connection ! Write(storage(0), Ack(storageOffset)) }
private def writeAll(): Unit = { for ((data, i) ← storage.zipWithIndex) { connection ! Write(data, Ack(storageOffset + i)) } }
//#storage-omitted } //#echo-handler
//#simple-echo-handler class SimpleEchoHandler(connection: ActorRef, remote: InetSocketAddress) extends Actor with ActorLogging {
import Tcp._
// sign death pact: this actor terminates when connection breaks context watch connection
case object Ack extends Event
def receive = { case Received(data) ⇒ buffer(data) connection ! Write(data, Ack)
//닫기를 시도해본다
context.become({ case Received(data) ⇒ buffer(data) case Ack ⇒ acknowledge() case PeerClosed ⇒ closing = true }, discardOld = false)
case PeerClosed ⇒ context stop self }
//#storage-omitted override def postStop(): Unit = { log.info(s”transferred $transferred bytes from/to [$remote]”) }
var storage = Vector.empty[ByteString] var stored = 0L var transferred = 0L var closing = false
val maxStored = 100000000L val highWatermark = maxStored * 5 / 10 val lowWatermark = maxStored * 3 / 10 var suspended = false
//#simple-helpers private def buffer(data: ByteString): Unit = { storage :+= data stored += data.size
if (stored > maxStored) { log.warning(s”drop connection to [$remote] (buffer overrun)”) context stop self
} else if (stored > highWatermark) { log.debug(s”suspending reading”) connection ! SuspendReading suspended = true } }
private def acknowledge(): Unit = { require(storage.nonEmpty, “storage was empty”)
val size = storage(0).size stored -= size transferred += size
storage = storage drop 1
if (suspended && stored < lowWatermark) { log.debug(“resuming reading”) connection ! ResumeReading suspended = false }
if (storage.isEmpty) { if (closing) context stop self else context.unbecome() } else connection ! Write(storage(0), Ack) } //#simple-helpers //#storage-omitted } //#simple-echo-handler