【问题标题】:Handling multiple TCP connections with Akka Actors使用 Akka Actor 处理多个 TCP 连接
【发布时间】:2015-06-01 07:27:39
【问题描述】:

我正在尝试使用akka Actor 设置一个简单的 TCP 服务器,它应该允许同时连接多个客户端。我将问题简化为以下简单程序:

package actorfail
import akka.actor._, akka.io._, akka.util._
import scala.collection.mutable._
import java.net._

case class Foo()

class ConnHandler(conn: ActorRef) extends Actor {
  def receive = {
    case Foo() => conn ! Tcp.Write(ByteString("foo\n"))
  }
}

class Server(conns: ArrayBuffer[ActorRef]) extends Actor {
  import context.system
  println("Listing on 127.0.0.1:9191")
  IO(Tcp) ! Tcp.Bind(self, new InetSocketAddress("127.0.0.1", 9191))
  def receive = {
    case Tcp.Connected(remote, local) =>
      val handler = context.actorOf(Props(new ConnHandler(sender)))
      sender ! Tcp.Register(handler)
      conns.append(handler)
  }
}

object Main {
  def main(args: Array[String]) {
    implicit val system = ActorSystem("Test")
    val conns = new ArrayBuffer[ActorRef]()
    val server = system.actorOf(Props(new Server(conns)))
    while (true)  {
      println(s"Sending some foos")
      for (c <- conns) c ! Foo()
      Thread.sleep(1000)
    }
  }
}

它绑定到 localhost:9191 并接受多个连接,将连接处理程序添加到全局数组并定期将字符串 "foo" 发送到每个连接。现在,当我尝试同时连接多个客户端时,只有第一个客户端获得“foo”。当我打开第二个连接时,它不会发送任何 foo,而是收到以下类型的日志消息:

Sending some foos
[INFO] [03/27/2015 21:24:07.331] [Test-akka.actor.default-dispatcher-6] [akka://Test/deadLetters] Message [akka.io.Tcp$Write] from Actor[akka://Test/user/$a/$b#-308726290] to Actor[akka://Test/deadLetters] was not delivered. [7] dead letters encountered. This logging can be turned off or adjusted with configuration settings 'akka.log-dead-letters' and 'akka.log-dead-letters-during-shutdown'.

我了解这意味着我们尝试向其发送Tcp.Write 命令的目标参与者不再接受消息。但这是为什么呢?你能帮我理解根本问题吗?我怎样才能做到这一点?

【问题讨论】:

  • 演员内的可变状态很好,但看起来您正试图从演员外部访问服务器演员的连接。这不是很Akka-y。为什么不让服务器管理连接池并在这些连接上发送消息?
  • @Gangstead 我刚刚测试了使用演员发送 Foos 并且这似乎有效,您能否详细说明此技术方面或记录在哪里?我希望实际上有一种方法可以将参与者与其他并发模型混合在一起。实际上,从外部访问参与者的代码是以阻塞方式从Source 读取行的循环。当然,我可以将其转换为演员,但将其表述为状态机会有点笨重。我宁愿有一种从外部发送消息的方法,只给 ActorRef。
  • 你从 Akka 项目的负责人那里得到了答案。他提供了一些文档链接。我认为您不能像希望的那样混合并发模型。
  • @Gangstead 好吧,根据他的回答我可以:)

标签: scala tcp akka actor


【解决方案1】:

上面的代码有两个问题:

  • 在参与者消息中发送可变状态并以非线程安全的方式对其进行变异
  • 在 Props 中包含不稳定的引用

在我详细说明之前,请考虑阅读文档,herehere,这些都在那里。

可变消息

ArrayBuffer 不是线程安全的,但是你将它从主程序传递给不同的参与者,然后他们独立(同时)修改它。这将导致更新丢失或数据结构本身损坏。另一方面,如果没有适当的同步,就不能保证主线程会看到修改,因为编译器原则上可以确定缓冲区在while 循环内没有变化并相应地优化代码。

actor 不依赖于共享的可变状态,而是只发送消息。在这种情况下,解决方案是将while 循环提升为参与者(但在一秒钟后将消息调度到self,而不是阻塞Thread.sleep(1000) 调用)。然后,连接处理程序只需要为这个foo 发送者actor 传递ActorRef,它们将向它发送一条消息以注册自己,然后该actor 将活动连接列表保留在其封装范围内。这样做的好处是您可以使用 DeathWatch 在连接终止时删除它们。

Props 中的不稳定引用

有问题的代码:Props(new ConnHandler(sender))

Props 是从一个 Actor 工厂构造的,在这种情况下,它被作为一个别名参数;整个new 表达式在稍后的时间被评估,每当这样一个actor被初始化时——可能在不同的线程上。这意味着sender 也会在稍后从执行上下文中评估,因此它可能是deadLetters(如果父actor当前没有运行,如果是,sender 可能会指向错误演员)。

这里的解决方案记录在here

【讨论】:

  • 我只在一个位置修改了 ArrayBuffer,我的印象是一个演员一次只能处理一条消息?当然它仍然是坏的,因为还有另一个演员在读取数组,但这只是为了演示目的,在真正的程序中并不是这样。不过,您发现了更有趣的问题,我稍后会尝试解决方案:)
  • 即使你只在一个位置修改了数组(你是对的,在这种特殊的、狭窄的情况下实际上是安全的),外部观察者可能会看到缓冲区处于不一致的状态,甚至可能会“倒退” ”演员在线程和 CPU 内核周围弹跳。共享可变状态不仅仅是明显的一类问题;-)
  • 这只是一个简单的例子,我在实际程序中不会做任何这样的事情。关闭构造函数参数上的道具有效,谢谢!
猜你喜欢
  • 2014-09-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-04
  • 2014-05-31
  • 1970-01-01
  • 2023-03-19
  • 2016-05-28
相关资源
最近更新 更多