【问题标题】:An Actor "queue"?演员“队列”?
【发布时间】:2010-06-07 22:03:00
【问题描述】:

在 Java 中,要编写一个向服务器发出请求的库,我通常会实现某种调度程序(与 Twitter4J 库中的调度程序不同:http://github.com/yusuke/twitter4j/blob/master/twitter4j-core/src/main/java/twitter4j/internal/async/DispatcherImpl.java)来限制连接数,以执行异步任务等

这个想法是创建 N 个线程。一个“任务”排队并通知所有线程,其中一个线程准备好后,将从队列中弹出一个项目,完成工作,然后返回等待状态。如果所有线程都在忙于一个任务,那么这个任务只是排队,下一个可用的线程将接受它。

这保持了与 N 的最大连接数,并允许最多 N 个任务同时运行。

我想知道我可以用 Actor 创建什么样的系统来完成同样的事情?有没有办法拥有 N 个 Actor,当一条新消息准备好时,将它传递给一个 Actor 来处理它 - 如果所有 Actor 都忙,就将消息排队?

【问题讨论】:

  • 您所描述的是一个线程池,因为 Java 5 它在标准库中,请参阅包 java.util.concurrent(类 ThreadPoolExecutor)。

标签: scala


【解决方案1】:

Akka Framework 旨在解决此类问题,并且正是您所寻找的。

看看这个docu - 有很多高度可配置的 dispather(基于事件、基于线程、负载平衡、工作窃取等)来管理参与者邮箱,并允许它们一起工作。您可能还会发现有趣的this blog post

例如。此代码基于固定线程池实例化新的 Work Stealing Dispatcher,实现其监管的 Actor 之间的负载平衡:

  val workStealingDispatcher = Dispatchers.newExecutorBasedEventDrivenWorkStealingDispatcher("pooled-dispatcher")
  workStealingDispatcher
  .withNewThreadPoolWithLinkedBlockingQueueWithUnboundedCapacity
  .setCorePoolSize(16)
  .buildThreadPool

使用调度器的Actor:

class MyActor extends Actor {

    messageDispatcher = workStealingDispatcher

    def receive = {
      case _ =>
    }
  }

现在,如果您启动 2 个以上的 actor 实例,调度程序将平衡 actor 的邮箱(队列)之间的负载(邮箱中有太多消息的 actor 将“捐赠”一些给没有任何消息的 actor做)。

【讨论】:

  • 我认为你可以对覆盖调度程序方法的 Scala 演员做同样的事情
【解决方案2】:

好吧,你必须了解演员调度程序,因为演员通常不是一对一的线程。演员背后的想法是你可能有很多,但实际的线程数将被限制在合理的范围内。他们也不应该长时间运行,而是快速回复他们收到的消息。简而言之,该代码的架构似乎与设计参与者系统的方式完全不一致。

尽管如此,每个工作的 Actor 可能会向 Queue Actor 发送一条消息,请求下一个任务,然后循环返回以做出反应。这个 Queue Actor 将接收队列消息或出队消息。可以这样设计:

val q: Queue[AnyRef] = new Queue[AnyRef]
loop {
  react {
    case Enqueue(d) => q enqueue d
    case Dequeue(a) if q.nonEmpty => a ! (q dequeue)
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-07-05
    • 2021-02-11
    • 2015-06-24
    • 2020-08-20
    • 1970-01-01
    • 1970-01-01
    • 2013-05-19
    • 2019-04-23
    相关资源
    最近更新 更多