【问题标题】:How to create dynamic consumer tasks for a Producer-Consumer solution using thread pools?如何使用线程池为生产者-消费者解决方案创建动态消费者任务?
【发布时间】:2020-02-14 03:04:32
【问题描述】:

我正在尝试实施生产者-消费者解决方案。

但我不想使用固定数量的消费者线程。相反,如果我的 eventQueue 已满,我想创建一个新的消费者线程。

我创建了一个 ExecutorService 但由于我只有一个 EventConsumerTask 实例,它只为这个任务创建了一个线程。

LinkedBlockingQueue<String> eventQueue = new LinkedBlockingQueue<>(50);

ExecutorService es = new ThreadPoolExecutor(5, 20, 60L, TimeUnit.SECONDS, new SynchronousQueue<Runnable>());

es.execute(new EventConsumerTask(eventQueue));

这是我的 EventConsumerTask;

public class EventConsumerTask implements Runnable{

    private LinkedBlockingQueue<String> eventQueue;

    public EventConsumerTask(LinkedBlockingQueue<String> eventQueue){
        this.eventQueue = eventQueue;
    }

    @Override
    public void run() {
        while(true) {
            try {
                String event = eventQueue.take();
                System.out.println(event);
                Thread.sleep(1000);
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }
}

【问题讨论】:

  • 你考虑过cached Thread Pool吗?
  • 与xingbin的回答相同的问题,不会添加更多EventConsumerTasks
  • 哦,我误解了这个概念。不过,这很愚蠢,为什么要在达到阈值后才开始并发处理任务,而不是让一组线程从单个队列中一个一个地处理单个事件?
  • 这可能是相当合理的,但前提是您可以放大
  • 不,我没有投反对票。

标签: java multithreading


【解决方案1】:

如果我理解正确,您想在队列已满时生成另一个 EventConsumerTask 吗?我想,会有一种方法可以填充队列,所以你可以在那里触发它:

    public synchronized void add(String item){
    if(eventQueue.remainingCapacity() == 0)
    {
        es.execute(new EventConsumerTask(eventQueue)); 
    }
    waitUntilQueueHasCapacityAgain();
    eventQueue.add(item);
}

如果您存储 EventConsumerTask,例如在列表中,如果队列变空,您可以再次缩减执行程序/线程。

顺便说一句。

    1234563 /em>
  • ConcurrentLinkedQueue 将在没有锁定的情况下执行类似的工作,你会想睡觉,机器人只有在没有可用的项目时,而不是在每个项目之后。

【讨论】:

  • 老实说,@danui 的答案对我来说看起来更规范,如果 Thread.sleep(1000); 不存在的话。现在我会寻找一些反应式解决方案,但这意味着您有一些异步接口来提供您的数据。
  • reactive streams查看http处理,他们从java 9开始就在那里。
【解决方案2】:

这并不能完全回答您的问题,但我建议采用这种一般模式:

// single thread to take from the event and dispatch handling to the pool
ExecutorService submitter = Executors.newSingleThreadExecutor();
// 5 - 20 threads for the individual event handlers, as suggested by xingbin
ExecutorService pooledExecutor = new ThreadPoolExecutor (5, 20, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>());

submitter.submit(() -> {
    while (true) {
        String event = eventQueue.take();
        // submit handling of the individual events in the pool
        pooledExecutor.submit(() -> {
            System.out.println(event);
            Thread.sleep(1000);
        }
    }
});

【讨论】:

    【解决方案3】:

    ThreadPoolExecutor 是动态的。

    查看doc

    public ThreadPoolExecutor​(int corePoolSize,
                              int maximumPoolSize,
                              long keepAliveTime,
                              TimeUnit unit,
                              BlockingQueue<Runnable> workQueue)
    

    corePoolSize - 保留在池中的线​​程数,即使它们 处于空闲状态,除非设置了 allowCoreThreadTimeOut

    maximumPoolSize - 的 池中允许的最大线程数

    new ThreadPoolExecutor(5, 20, 60L, TimeUnit.SECONDS, new SynchronousQueue<Runnable>());
    

    这意味着线程池最多可以容纳60个线程。

    【讨论】:

    • 我不认为,它回答了他的要求,它不会添加更多的 EventConsumerTasks
    猜你喜欢
    • 2016-08-08
    • 1970-01-01
    • 2013-11-10
    • 1970-01-01
    • 1970-01-01
    • 2018-09-24
    • 2017-02-01
    • 2012-04-30
    • 1970-01-01
    相关资源
    最近更新 更多