【问题标题】:Processing sub-streams of a stream in Java using executors使用执行器在 Java 中处理流的子流
【发布时间】:2015-11-30 06:41:55
【问题描述】:

我有一个程序可以处理通过网络传入的大量数据流(不是java.util.stream,而是InputStream)。流由对象组成,每个对象都有一种子流标识符。现在整个处理都是在一个线程中完成的,但是需要大量的CPU时间,并且每个子流都可以很容易地独立处理,所以我正在考虑多线程。

但是,每个子流都需要保持大量庞大的状态,包括各种缓冲区、哈希映射等。没有特别的理由让它并发或同步,因为子流彼此独立。此外,每个子流都要求其对象按照它们到达的顺序进行处理,这意味着每个子流可能应该有一个线程(但可能一个线程处理多个子流)。

我正在考虑几种方法,但它们不是很优雅。

  1. 为所有任务创建一个 ThreadPoolExecutor。每个任务将包含下一个要处理的对象以及对保持所有状态的Processor 实例的引用。这将确保必要的发生前关系,从而确保处理线程将看到该子流的最新状态。据我所知,这种方法无法确保同一子流的下一个对象将在同一线程中处理。此外,它需要保证对象将按照它们进入的顺序进行处理,这将需要额外同步Processor 对象,从而引入不必要的延迟。

  2. 手动创建多个单线程执行器和一种将子流标识符映射到执行器的哈希映射。这种方法需要手动管理执行器,在新的子流开始或结束时创建或关闭它们,并相应地在它们之间分配任务。

  3. 创建一个自定义执行器来处理一个特殊的任务子类,每个任务都有一个子流 ID。该执行器将使用它作为提示,以使用与前一个具有相同 ID 的线程执行此任务的相同线程。但是,我没有看到实现这种执行器的简单方法。不幸的是,似乎无法扩展任何现有的执行程序类,从头开始实现执行程序有点过头了。

  4. 创建一个ThreadPoolExecutor,但不是为每个传入对象创建一个任务,而是为每个子流创建一个长时间运行的任务,该任务将阻塞在并发队列中,等待下一个对象。然后根据对象的子流 ID 将对象放入队列中。这种方法需要与子流一样多的线程,因为任务将被阻塞。子流的预期数量约为 30-60,因此可以接受。

  5. 或者,按照 4 进行,但限制线程数,将多个子流分配给单个任务。这是 2 和 4 之间的一种混合。据我所知,这是其中最好的方法,但它仍然需要在任务之间进行某种手动子流分配,以及关闭额外任务的某种方法子流结束。

确保每个子流在其自己的线程中处理而没有大量容易出错的代码的最佳方法是什么?这样下面的伪代码就可以工作了:

// loop {
    Item next = stream.read();
    int id = next.getSubstreamID();
    Processor processor = getProcessor(id);
    SubstreamTask task = new SubstreamTask(processor, next, id);
    executor.submit(task); // This makes sure that the task will
                           // be executed in the same thread as the
                           // previous task with the same ID.
// } // loop

【问题讨论】:

    标签: java multithreading concurrency java.util.concurrent threadpoolexecutor


    【解决方案1】:

    我建议使用一组单线程执行器。如果您可以为子流设计一致的散列策略,则可以将子流映射到各个线程。例如

    final ExecutorsService[] es = ...
    
    public void submit(int id, Runnable run) {
       es[(id & 0x7FFFFFFF) % es.length].submit(run);
    }
    

    密钥可以是String 或long,但可以通过某种方式识别子流。如果您知道某个特定的子流非常昂贵,则可以为其分配一个专用线程。

    【讨论】:

    • 流同样昂贵,所以这应该不是问题。而且我只是保持固定数量的执行器处于活动状态,而不是关闭它们并根据活动子流的数量创建新的执行器?所以这是我方法 2 的简化,对吧?
    • @SergeyTachenov 正确,如果执行者无事可做,他们不会浪费 CPU。如果您知道不再需要它们,我会关闭它们。
    • 这种方法的一个问题是,无论我选择什么散列策略,它都不是最优的。有时我会发现一些执行者做了很多工作,而另一些则什么都不做。所以我想我必须引入某种从 ID 到执行器的动态映射,每次出现新的子流时都要修改它,保留现有子流的映射。
    • @SergeyTachenov 不要忘记您的线程没有固定到给定的 CPU。如果每个逻辑 CPU 有一个线程,那么您很可能有足够的线程来保持每个核心忙碌,并且如果您说线程数量是逻辑 CPU 的两倍,那么您很可能会让它们都忙碌。您可以调整线程的数量来品尝。在极端情况下,每个子流都有一个线程。
    • 最后我使用了一个ConcurrentHashMap,它将一个ID映射到执行器索引。如果没有分配执行器,我会分配一个新的执行器,从那些已经分配了最少子流的执行器中进行选择。使用computeIfAbsent 效果很好。当分配给一个执行器的多个子流完成时,我仍然会得到不均匀的分布,但我通过以某种方式分配它们来缓解它,以便将可能同时完成的那些分配给不同的执行器。
    【解决方案2】:

    我最终选择的方案是这样的:

    private final Executor[] streamThreads
            = new Executor[Runtime.getRuntime().availableProcessors()];
    {
        for (int i = 0; i < streamThreads.length; ++i) {
            streamThreads[i] = Executors.newSingleThreadExecutor();
        }
    }
    private final ConcurrentHashMap<SubstreamId, Integer>
            threadById = new ConcurrentHashMap<>();
    

    此代码确定使用哪个执行器:

        Message msg = in.readNext();
        SubstreamId msgSubstream = msg.getSubstreamId();
        int exe = threadById.computeIfAbsent(msgSubstream,
                id -> findBestExecutor());
        streamThreads[exe].execute(() -> {
            // processing goes here
        });
    

    findBestExecutor() 函数是这样的:

    private int findBestExecutor() {
        // Thread index -> substream count mapping:
        final int[] loads = new int[streamThreads.length];
        for (int thread : threadById.values()) {
            ++loads[thread];
        }
        // return the index of the minimum load
        return IntStream.range(0, streamThreads.length)
                .reduce((i, j) -> loads[i] <= loads[j] ? i : j)
                .orElse(0);
    }
    

    当然,这不是很有效,但请注意,此函数仅在出现新的子流时调用(每隔几个小时会发生几次,所以在我的情况下这没什么大不了的)。我的真实代码看起来有点复杂,因为我有办法确定两个子流是否可能同时完成,如果是,我会尝试将它们分配给不同的线程,以便在它们完成后保持均匀的负载。但由于我从未在问题中提到过这个细节,我想它也不属于答案。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-02-24
      • 2020-07-14
      • 2017-06-26
      • 2016-05-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多