【问题标题】:Actors, ForkJoinPool, and ordering of messagesActors、ForkJoinPool 和消息排序
【发布时间】:2022-01-18 17:21:34
【问题描述】:

我需要帮助了解 Actor 系统如何使用 ForkJoinPool 并维护排序保证。

我一直在玩 Actr https://github.com/zakgof/actr,这是一个简单的小型演员系统。我认为我的问题也适用于 Akka。我有一段简单的代码可以发送一个演员号码 1 到 10。演员只是打印消息;并且消息不按顺序排列。我得到 1,2,4,3,5,6,8,7,9,10。

我认为这与 ForkJoinPool 有关。 Actr 将消息包装到 Runnable 中并将其发送到 ForkJoin Executor。当任务执行时,它会将消息放到目标 Actor 的队列中并对其进行处理。我对 ForkJoinPool 的理解是任务被分配到多个线程。我添加了日志记录,消息 1,2,3,... 被分发到不同的线程,消息被乱序放入 Actor 的队列中。

我错过了什么吗? Actr 的 Scheduler 类似于 Akka 的 Disapatcher,可以在这里找到:https://github.com/zakgof/actr/blob/master/src/main/java/com/zakgof/actr/impl/ExecutorBasedScheduler.java

ExecutorBasedScheduler 使用 ForkJoinPool.commonPool 构造,如下所示:

public static IActorScheduler newForkJoinPoolScheduler(int throughput) {
    return new ExecutorBasedScheduler(ForkJoinPool.commonPool(), throughput);
}

Actor 如何使用 ForkJoinPool 并保持消息有序?

【问题讨论】:

    标签: akka actor forkjoinpool


    【解决方案1】:

    我根本无法与 Actr 交谈,但在 Akka 中,单个消息不会创建为 ForkJoinPool 任务。 (由于许多原因,每条消息一个任务似乎是一种非常糟糕的方法,而不仅仅是排序问题。也就是说,通常可以非常快速地处理消息,如果每条消息有一个任务,开销会非常高。你想要一些批处理,至少在负载下,这样您就可以获得更好的线程局部性和更少的开销。)

    本质上,在 Akka 中,actor 邮箱是对象中的队列。当邮箱收到一条消息时,它会检查它是否已经安排了一个任务,如果没有,它将向 ForkJoinPool 添加一个新任务。所以 ForkJoinPool 任务不是“处理这条消息”,而是“处理与这个特定 Actor 的邮箱关联的 Runnable”。在任务被安排并且 Runnable 运行之前,显然会经过一段时间。当 Runnable 运行时,邮箱可能已收到更多消息。但它们将刚刚被添加到队列中,然后 Runnable 将按照配置的顺序处理尽可能多的它们,按照接收它们的顺序。

    这就是为什么在 Akka 中,你可以保证一个邮箱中消息的顺序,但不能保证发送给不同 Actor 的消息的顺序。如果我将消息 A 发送给 Actor Alpha,然后将消息 B 发送给 Actor Beta,然后将消息 C 发送给 Actor Alpha,我可以保证 A 将在 C 之前。但是 B 可能发生在 A 和 C 之前、之后或同时.(因为 A 和 C 将由同一个任务处理,但 B 将是不同的任务。)

    Messaging Ordering Docs :关于什么是有保证的,什么不是关于订购的更多细节。

    Dispatcher Docs:调度程序是 Actor 和实际执行之间的连接。 ForkJoinPool 只是一种实现(尽管是一种非常常见的实现)。

    编辑:只是想我会添加一些指向 Akka 源的链接来说明。请注意,这些都是内部 API。 tell 是你如何使用它,这一切都在幕后。 (我正在使用永久链接,这样我的链接就不会被破坏,但请注意 Akka 可能在您使用的版本中发生了变化。)

    密钥位在akka.dispatch.Dispatcher.scala

    您的tell 将经过一些步骤才能到达正确的邮箱。但最终:

    • 调用dispatch 方法将其加入队列。这个很简单,入队并调用registerForExecution方法
    • registerForExecution 这个方法实际上是先检查是否需要调度。如果它需要调度它使用 executorService 来调度它。请注意,executorService 是抽象的,但在提供邮箱作为参数的服务上调用 execute
    • execute 如果我们假设实现是 ForkJoinPool,这就是我们最终使用的 executorService 执行方法。本质上,我们只需使用提供的参数(邮箱)作为可运行对象创建一个 ForkJoinTask。
    • run Mailbox 方便地是 Runnable,因此 ForkJoinPool 最终将在调度后调用此方法。您可以看到它处理特殊的系统消息,然后调用processMailbox,然后(最终)再次调用registerForExecution。请注意,registerForExecution 首先检查它是否需要调度,所以这不是一个无限循环,它只是检查是否还有剩余工作要做。当我们在 Mailbox 类中时,您还可以查看我们在 Dispatcher 中使用的一些方法,以查看是否需要调度、将消息实际添加到队列等。
    • processMailbox 本质上只是一个调用 actor.invoke 的循环,除了它必须做大量检查以查看它是否有系统消息、是否停止工作、是否超过阈值、是否已被中断等。
    • invoke 是您编写的代码(receiveMessage)实际被调用的地方。

    如果你真的点击了所有这些链接,你会发现我简化了很多。有很多错误处理和代码来确保一切都是线程安全的、超级高效的和防弹的。但这就是代码流的要点。

    【讨论】:

    • 让我们看看我是否明白你在说什么。 1:我调用 printActor.tell("One")。消息进入 printActor 的邮箱队列。 2:如果 mainbox 已经从空转换为非空,则将任务放入 ForkJoinPool。 3:任务将排空并处理队列上的消息,直到队列为空(或有阈值逻辑,但这是一个旁白)。 4:所以在任何时候都应该只有一个线程处理特定Actor的队列。我说的对吗?
    • 是的。需要注意的是,ForkJoinPool 只是一个调度程序实现。还有其他调度程序实现,例如 PinnedDispatcher,其中有一个专用线程。 (我对“应该”这个词也有点退缩,因为它是有保证的,它是 Akka 的核心原则。但我想你明白了。)
    • 添加一些代码链接只是为了说明。
    猜你喜欢
    • 2021-06-19
    • 1970-01-01
    • 2016-01-15
    • 2015-05-27
    • 2012-10-06
    • 2015-12-21
    • 2019-11-19
    • 2017-06-24
    • 2015-11-27
    相关资源
    最近更新 更多