【问题标题】:Processing a large number of tasks with CompletionService使用 CompletionService 处理大量任务
【发布时间】:2014-05-29 03:51:21
【问题描述】:

我需要在多核机器上处理大量(>1 亿)请求(每个请求是处理数据文件中的一行,并且涉及到与远程系统的一些 I/O。虽然细节确实没关系,具体任务是从一些数据文件中加载分布式 Hazelcast 地图)。执行将通过ThreadPoolExecutor 处理。一个线程将读取文件,然后将数据提交给多个独立线程以将其放入映射中。机器有 32 个核心,所以有足够的可用于并行加载地图。

由于请求数量众多,一般的创建任务并排队到执行器服务的做法是行不通的,因为排队的任务会占用太多的内存。

这带来了ExecutorCompletionService。有了它,将在先前的操作完成时提交任务,这可以通过调用take()(或poll(),如适用)得知。当执行器服务的所有线程都使用时,这将正常工作。但是,“加载所有线程”还没有完成。有两个阶段:

  • 填满队列:当池中仍有未使用的线程时,将任务提交给 ExecutorCompletionService,不要等到提交更多

  • 喂入队列:一旦线程全部使用完毕,只有在前一个任务完成后才提交任务。因此,将尽可能快地提供行,但不会更快,也不会排队。

上面可以编码,但我想知道上面的逻辑是否已经实现,我不知何故错过了它。我问是因为它看起来很常见。

【问题讨论】:

    标签: java multithreading executorservice hazelcast completion-service


    【解决方案1】:

    您可以在创建ThreadPoolExecutor 时指定BlockingQueue 实现。如果您要避免的只是创建过多的Runnable 对象,那么您可以使用有界BlockingQueue,例如ArrayBlockingQueue 有一个线程将项目推送到队列中,当队列满负荷时,该线程将被阻塞。

    【讨论】:

    • 我正在尝试blockingQueue,但它对Callable没有帮助。
    【解决方案2】:

    如果我理解您的要求,(如果我错了,请纠正我)然后您需要一种机制,其中有多个任务,并且您需要最多 n 任务并行执行,其他任务应该在队列中等待,但是一旦你提交了一个任务,你就不想闲逛或让线程提交任务忙,它可以继续它的工作

    对于同样的场景,我们混合使用LinkedBlockingQueueThread,相信一个简单的函数可以帮助你理解,

    private final LinkedBlockingQueue<YourTaskObjType> EnqueuedTasks;
    
    private void initTasksProcessingThreads(int numberOfThreads) 
    {
        EnqueuedTasks= new LinkedBlockingQueue<YourTaskObjType>();
        for (int i = 0; i < numberOfThreads; i++) 
        {
            // each thread will run forever and process incoming
            //Change requests
            Thread worker = new Thread(new Runnable() 
            {               
                public void run() 
                {
                    while (true) 
                    {
                        try 
                        {   
                            YourTaskObjType task = EnqueuedTasks.take(); //This will wait infinitely until tasks are available
                            PerformTask(task); //Your function which will perform the task operation
                        } 
                        catch (InterruptedException e) 
                        {                                                                
                            Thread.currentThread().interrupt();
                            return;
                        } 
                        catch(Exception e)
                        {
                            e.printStackTrace();
                        }
                    }
                }
            });         
            worker.start();
        }
    }
    

    然后你可以使用一个简单的函数将Tasks添加到LinkedBlockingQueue

    public void AddTask(YourTaskObjType TaskObj)
    {
        EnqueuedTasks.put(TaskObj);                         
    }
    

    【讨论】:

    • 看起来您正在尝试重塑 ThreadPoolExecutor,这有什么更好的地方?
    • @SimonC Reinvent ThreadPoolExecutor ???????????????这是我做的方式....而且我从来没有提到它更好....-1因为没有实现你认为的方式.....
    • 您的代码执行与ThreadPoolExecutor 相同的工作,只是可靠性和可维护性要低得多。使用ThreadPoolExecutor 并没有增加任何内容。 -1 是为了在没有任何充分理由的情况下尝试在标准库中重新实现一个类。
    猜你喜欢
    • 1970-01-01
    • 2021-11-17
    • 2018-03-05
    • 1970-01-01
    • 2021-06-29
    • 2022-01-18
    • 1970-01-01
    • 2013-08-23
    • 2014-12-20
    相关资源
    最近更新 更多