【问题标题】:Designing Akka Supervisor Hierarchy设计 Akka 主管层次结构
【发布时间】:2015-06-16 01:28:11
【问题描述】:

请注意:我是一名 Java 开发人员,没有 Scala 的工作知识(很遗憾)。我会要求答案中提供的任何代码示例都将使用 Akka 的 Java API。

我对 Akka 和演员来说是全新的,并且正在尝试建立一个相当简单的演员系统:

因此,DataSplitter actor 运行并将相当大的二进制数据块(例如 20GB)拆分为 100KB 的块。对于每个块,数据通过DataCacher 存储在DataCache 中。在后台,DataCacheCleaner 在缓存中翻找并找到可以安全删除的数据块。这就是我们防止缓存变为 20GB 大小的方法。

在将块发送到DataCacher 进行缓存后,DataSplitter 然后通知ProcessorPool 现在需要处理的块。 ProcessorPool 是一个路由器/池,由 数万 个不同的 ProcessorActors 组成。当每个ProcessActor 收到“处理” 100KB 数据块的通知时,它会从DataCacher 获取数据并对其进行一些处理。

如果你想知道我为什么还要在这里缓存任何东西(因此有 DataCacherDataCacheDataCacheCleaner),我的想法是 100KB 仍然是一个相当大的消息传递给数以万计的演员实例(100KB * 1,000 = 100MB),所以我试图只存储一次 100KB 块(在缓存中),然后让每个演员通过缓存 API 通过引用访问它。

还有一个Mailmanactor订阅了事件总线,拦截了所有DeadLetters

所以,总共有 6 个演员:

  • DataSplitter
  • DataCacher
  • DataCacheCleaner
  • ProcessorPool
  • ProcessorActor
  • Mailman

Akka 文档宣扬您应该根据划分子任务而不是纯粹按功能来分解您的 Actor 系统,但我并不完全了解这在此处如何应用。手头的问题是我正在尝试在这些参与者之间组织一个主管层次结构,我不确定最好/正确的方法是什么。显然ProcessorPool 是一个需要成为ProcessorActors 的父/主管的路由器,所以我们有这个已知的层次结构:

/user/processorPool/
    processorActors

但除了已知/明显的关系之外,我不确定如何组织我的其他演员。我可以让他们在一个共同/主要演员下成为“同行”:

/user/master/
    dataSplitter/
    dataCacher/
    dataCacheCleaner/
    processorPool/
        processorActors/
    mailman/

或者我可以省略 master(根)actor 并尝试使缓存周围的东西更加垂直:

/user/
    dataSplitter/
    cacheSupervisor/
        dataCacher/
        dataCacheCleaner/
    processorPool/
        processorActors/
    mailman/

对 Akka 如此陌生,我只是不确定最好的行动方案是什么,如果有人可以在这里帮助一些初步的手把手,我相信灯泡都会打开。而且,就像组织这个层次结构一样重要,我什至不确定我可以使用什么 API 构造来在代码中实际创建层次结构

【问题讨论】:

  • 你最后是怎么设计这个的?

标签: java akka actor fault-tolerance akka-supervision


【解决方案1】:

将它们组织在一个master 下更易于管理,因为您可以通过主管访问所有演员watched(在本例中为master)。

一个分层实现可以是:

主监演员

class MasterSupervisor extends UntypedActor {

private static SupervisorStrategy strategy = new AllForOneStrategy(2,
        Duration.create(5, TimeUnit.MINUTES),

        new Function<Throwable, Directive>() {
            @Override
            public Directive apply(Throwable t) {

                if (t instanceof SQLException) {
                    log.error("Error: SQLException")
                    return restart()
                } else if (t instanceof IllegalArgumentException) {
                    log.error("Error: IllegalArgumentException")
                    return stop()
                } else {
                    log.error("Error: GeneralException")
                    return stop()
                }
            }
        });

@Override
public SupervisorStrategy supervisorStrategy() { return strategy }

@Override
void onReceive(Object message) throws Exception {
     if (message.equals("SPLIT")) {
          // CREATE A CHILD OF MyOtherSupervisor
          if (!dataSplitter) {
              dataSplitter = context().actorOf(FromConfig.getInstance().props(Props.create(DataSplitter.class)), "DataSplitter")

              // WATCH THE CHILD
              context().watch(dataSplitter)

              log.info("${self().path()} has created, watching and sent JobId = ${message} message to DataSplitter")
          }

          // do something with message such as Forward
          dataSplitter.forward(message, context())
      }
}

DataSplitter Actor

class DataSplitter extends UntypedActor {

    // Inject a Service to do the main operation
    DataSplitterService dataSplitterService

    @Override
    void onReceive(Object message) throws Exception {
        if (message.equals("SPLIT")) {
            log.info("${self().path()} recieved message: ${message} from ${sender()}")
            // do something with message such as Forward
            dataSplitterService.splitData()
        }
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-11-21
    • 2011-05-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-03-27
    • 2016-04-22
    • 2014-12-06
    相关资源
    最近更新 更多