【问题标题】:Is there a way to ensure that threads are assigned to a specified set of objects?有没有办法确保将线程分配给一组指定的对象?
【发布时间】:2012-11-01 03:57:07
【问题描述】:

我们正在开发一个应用程序,其中一组对象可以通过接收来自 3 个不同来源的消息而受到影响。每条消息(来自任何来源)都有一个对象作为其目标。每个消息接收器都将在自己的线程上运行。

我们希望消息的处理(接收后)尽可能高速,因此针对目标对象的消息处理将由线程池中的另一个线程完成。消息的处理将比阅读/接收来自发件人的消息花费更长的时间。

我认为如果池中的每个线程只专用于一组特定的对象会更快,例如:

Thread1 -> objects named A-L
Thread2 -> objects named M-Z

每组对象(或线程)都有一个专用的待处理消息队列。

我的假设是,如果唯一需要的线程同步是在每个接收线程和一个处理线程之间,在需要将消息放入阻塞队列的持续时间内,它会比随机分配工作线程更快处理消息(在这种情况下,可能有 2 个不同的线程处理同一对象的消息)。

我的问题实际上是两部分:

  1. 人们是否同意专用工作线程的假设 对一组特定的对象是更好/更快的方法吗?

  2. 假设这是一种更好的方法,请执行现有的 Java ThreadPool 类有办法支持这一点吗?还是需要我们编码 我们自己的线程池实现?

感谢您提供的任何建议。

【问题讨论】:

  • 所以诀窍是不能让 2 条具有相同目标对象的消息同时被不同的线程处理?
  • 你有没有做过任何测试,看看这是否是一个问题?
  • @Gray:是的 - 同一目标对象的 2 条消息中的每一条都将更新目标对象的状态,并且还可能触发生成另一条消息。这两条消息绝对不允许同时更新对象(即没有同步),因为这可能导致对象的状态不一致(如并发线程所见)。
  • @MattiLyra:从其他类似应用程序的性能来看,我们知道接收延迟与处理延迟之间存在不平衡。关于并发线程是否会看到不一致的状态(如果没有按照我的建议进行同步或专用),我不能 100% 确定 - 但这是一个合理的假设。
  • 但是如果目标对象必须同步,那么划分对象集的意义何在?为什么要将对象分成 A-L M-Z?

标签: java multithreading threadpool


【解决方案1】:

作为替代方法:我建议为此使用现有框架,例如 RabbitMQActiveMQ。尝试发明自己的消息传递框架可能是一个挑战。如果您尝试使用自己的框架增加价值,那是一回事。如果您只需要一个来实现您的目标,那就是另一个。这些框架提出了许多优化消息传递的选项,值得考虑。

【讨论】:

  • 实际上,我们并没有发明自己的框架,我们正在与一些通过 TCP 发送消息的外部系统集成。我认为将 MQ 解决方案作为中间件层肯定会增加延迟和复杂性。
【解决方案2】:

我的答案是:

  • 1 - 是的
  • 2 -
    • a) 否
    • b) 你不需要

一些解释:

  • 您希望一项任务根据某种算法将消息分发到不同的队列,
  • 您希望每个消息队列有一个任务从其分配的队列中提取消息并处理它们。

我不认为这些前提与线程池的目的相矛盾,线程池只是将任务与线程相关联。不过,在这个模型中,线程池只会将线程与任务关联一次,然后线程会继续运行以轮询其输入消息队列。

线程的摩擦点应该是中间消息队列,可能还有与这些消息处理相关的其他资源。根据您的解释,我想您计划通过巧妙地将消息处理划分为任务来将第二种类型减少到最低限度。每个队列只能被与队列关联的分区任务和处理任务访问,因此应该是最小的。

【讨论】:

  • 能否请您指点我一些示例代码,或者哪个 Java ThreadPool/ExecutorService 可以很容易地做到这一点?
  • 设置中不会有什么特别困难的事情。只需将任务定义为无限循环轮询其队列,并对新消息采取行动。线程池会启动它们一次,然后让它们自己运行。
  • 我还不能为您提供示例应用程序(我不在我的电脑上)。
  • 如果你有n消息分区,那么你只需要n+1线程,所以一个fixedThreadPool就可以了。
【解决方案3】:

一般来说,这样的方法是个坏主意。它属于“don't optimize early”的口头禅。

此外,如果实施你的想法可能损害你的表现,而不是帮助它。一个无法正常工作的简单示例是,如果您突然在一种类型上收到大量请求 - 另一个工作线程将处于空闲状态。

最好的方法是使用标准的生产者-消费者模式,并通过在各种负载下进行系统测试来调整消费者线程的数量——最好是通过输入真实交易的记录。

这些情况的“转到”框架是来自java.util.concurrent 包的类。我建议使用BlockingQueue(可能是ArrayBlockingQueue)和从Executors 工厂方法之一创建的ExecutorService,可能是newCachedThreadPool()


实施并进行系统测试后,如果您发现已证实的性能问题,请分析您的系统,找到瓶颈并修复它。

你不应该及早优化的原因是大多数时候问题不在你期望的地方

【讨论】:

  • 我正在考虑如何以线程中立的方式构建它,以便可以更改线程算法。因此,如果我能想出一种方法来抽象某个接口后面的消息传递,那么也许我可以将这些决定留到我们进行性能测试之前。
  • Executors 类提供了几种工厂方法,其中一些是可配置的,而ExecutorService 本身是可配置的,因此您始终可以编写自己的实现。这是一个很好的起点,也是一个很好的学习 API。
【解决方案4】:

[Is] 将工作线程专用于一组特定的对象是一种更好/更快的方法吗?

我认为总体目标是尝试最大化这些入站消息的并发处理。您有来自 3 个来源的接收器,它们需要将消息放入将得到最佳处理的池中。因为来自 3 个源中的任何一个的消息都可能处理不能同时处理的同一目标对象,所以您希望以某种方式分割您的消息,以便可以同时处理它们,但前提是它们保证不会同时处理引用同一个目标对象。

我会在您的目标对象上实现hashCode() 方法(可能只是name.hashCode()),然后使用该值将对象放入BlockingQueues 的数组中,每个对象都有一个线程使用它们。使用 Executors.newSingleThreadExecutor() 的数组就可以了。通过队列数修改哈希值模式并将其放入该队列中。您需要将处理器的数量预先定义为最大值。取决于处理的 CPU 密集程度。

所以类似下面的代码应该可以工作:

 private static final int NUM_PROCESSING_QUEUES = 6;
 ...
 ExecutorService[] pools = new ExecutorService[NUM_PROCESSING_QUEUES];
 for (int i = 0; i < pools.length; i++) {
    pools[i] = Executors.newSingleThreadExecutor();
 }
 ...
 // receiver loop:
 while (true) {
    Message message = receiveMessage();
    int hash = Math.abs(message.hashCode());
    // put each message in the appropriate pool based on its hash
    // this assumes message is runnable
    pools[hash % pools.length].submit(message);
 }

这种机制的一个好处是您可以限制有关目标对象的同步。你知道同一个目标对象只会被一个线程更新。

人们是否同意将工作线程专用于一组特定对象是一种更好/更快的方法的假设?

是的。这似乎是获得最佳并发性的正确方法。

假设这是一种更好的方法,现有的 Java ThreadPool 类是否有办法支持这种方法?还是需要我们编写自己的 ThreadPool 实现?

我不知道有任何线程池可以做到这一点。但是,我不会编写您自己的实现。就像上面的代码大纲一样使用它们。

【讨论】:

  • 感谢您的示例代码。我会试一试,然后告诉你它是如何工作的。
  • 感谢您的回答。请记住,.hashCode() 可能会返回一个负整数。为了安全使用:Math.abs(hash).
  • 感谢@RubenErnst。我已经确定了答案。
【解决方案5】:

您应该能够为 ThreadPoolExecutor 提供一个特殊的 BlockingQueue。队列会记住哪个线程正在处理哪种类型的消息,以便它可以保留所有相同类型的消息。

MyQueue

    ownership relation of thread - msgType 

    take/poll()

        if current thread owns msg type X
            if there is a message of type X
                return that message
            else
                give up ownership

        // current thread does not own any message type
        if there is a messsage of type Y, Y is not owned by any thread
            current thread owns Y
            return that message

        // there's no message belonging to an unowned type
        wait then retry 

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-07-01
    • 2019-04-04
    • 2012-04-24
    • 2014-12-04
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多