Skip to content
dev-marco-song
Go back

maven으로 scala AKKA echo 서버 구현하기

Edit page

maven으로 scala AKKA echo 서버 구현하기

자바, 스칼라의 오픈소스 배포용으로 메이븐이 많이 사용되고 있다. 메이븐을 이용한 스칼라 프로젝트 사용법을 정리하고자 이 포스트를 작성한다.

소스 프로젝트는 깃허브로 배포해두었다. https://github.com/jamsya/akka-echo


메이븐에서 자동으로 생성하는 프로젝트로는 컴파일이 되지 않았다.

$ mvn archetype:generate

스칼라 버전 바꿔보고, 지웠다 깔고, jdk 바꿔보고 다 해봤는데 소용 없었다.

헌데 git에서 예제로 받은 코드는 메이븐으로 컴파일이 잘 된다 -_-;; 지금 메이븐에서의 스칼라 프로젝트 생성에는 문제가 있어보인다. 아니면 내가 모르는게 또 있거나..

그래서 깃헙에 있는 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


Edit page
Share this post on:

Previous Post
POCO C++ 라이브러리 에코 서버 구현하기 / cmake 기반
Next Post
리눅스 환경에서 git 사용법