【发布时间】:2015-11-30 06:41:55
【问题描述】:
我有一个程序可以处理通过网络传入的大量数据流(不是java.util.stream,而是InputStream)。流由对象组成,每个对象都有一种子流标识符。现在整个处理都是在一个线程中完成的,但是需要大量的CPU时间,并且每个子流都可以很容易地独立处理,所以我正在考虑多线程。
但是,每个子流都需要保持大量庞大的状态,包括各种缓冲区、哈希映射等。没有特别的理由让它并发或同步,因为子流彼此独立。此外,每个子流都要求其对象按照它们到达的顺序进行处理,这意味着每个子流可能应该有一个线程(但可能一个线程处理多个子流)。
我正在考虑几种方法,但它们不是很优雅。
为所有任务创建一个
ThreadPoolExecutor。每个任务将包含下一个要处理的对象以及对保持所有状态的Processor实例的引用。这将确保必要的发生前关系,从而确保处理线程将看到该子流的最新状态。据我所知,这种方法无法确保同一子流的下一个对象将在同一线程中处理。此外,它需要保证对象将按照它们进入的顺序进行处理,这将需要额外同步Processor对象,从而引入不必要的延迟。手动创建多个单线程执行器和一种将子流标识符映射到执行器的哈希映射。这种方法需要手动管理执行器,在新的子流开始或结束时创建或关闭它们,并相应地在它们之间分配任务。
创建一个自定义执行器来处理一个特殊的任务子类,每个任务都有一个子流 ID。该执行器将使用它作为提示,以使用与前一个具有相同 ID 的线程执行此任务的相同线程。但是,我没有看到实现这种执行器的简单方法。不幸的是,似乎无法扩展任何现有的执行程序类,从头开始实现执行程序有点过头了。
创建一个
ThreadPoolExecutor,但不是为每个传入对象创建一个任务,而是为每个子流创建一个长时间运行的任务,该任务将阻塞在并发队列中,等待下一个对象。然后根据对象的子流 ID 将对象放入队列中。这种方法需要与子流一样多的线程,因为任务将被阻塞。子流的预期数量约为 30-60,因此可以接受。或者,按照 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