`
wbj0110
  • 浏览: 1598056 次
  • 性别: Icon_minigender_1
  • 来自: 上海
文章分类
社区版块
存档分类
最新评论

ExecutorCompletionService

    博客分类:
  • Java
阅读更多

当我们通过Executor提交一组并发执行的任务,并且希望在每一个任务完成后能立即得到结果,有两种方式可以采取:

 

方式一:

通过一个list来保存一组future,然后在循环中轮训这组future,直到每个future都已完成。如果我们不希望出现因为排在前面的任务阻塞导致后面先完成的任务的结果没有及时获取的情况,那么在调用get方式时,需要将超时时间设置为0 

Java代码  收藏代码
  1. public class CompletionServiceTest {  
  2.   
  3.     static class Task implements Callable<String>{  
  4.         private int i;  
  5.           
  6.         public Task(int i){  
  7.             this.i = i;  
  8.         }  
  9.   
  10.         @Override  
  11.         public String call() throws Exception {  
  12.             Thread.sleep(10000);  
  13.             return Thread.currentThread().getName() + "执行完任务:" + i;  
  14.         }     
  15.     }  
  16.       
  17.     public static void main(String[] args){  
  18.         testUseFuture();  
  19.     }  
  20.       
  21.     private static void testUseFuture(){  
  22.         int numThread = 5;  
  23.         ExecutorService executor = Executors.newFixedThreadPool(numThread);  
  24.         List<Future<String>> futureList = new ArrayList<Future<String>>();  
  25.         for(int i = 0;i<numThread;i++ ){  
  26.             Future<String> future = executor.submit(new CompletionServiceTest.Task(i));  
  27.             futureList.add(future);  
  28.         }  
  29.                   
  30.         while(numThread > 0){  
  31.             for(Future<String> future : futureList){  
  32.                 String result = null;  
  33.                 try {  
  34.                     result = future.get(0, TimeUnit.SECONDS);  
  35.                 } catch (InterruptedException e) {  
  36.                     e.printStackTrace();  
  37.                 } catch (ExecutionException e) {  
  38.                     e.printStackTrace();  
  39.                 } catch (TimeoutException e) {  
  40.                     //超时异常直接忽略  
  41.                 }  
  42.                 if(null != result){  
  43.                     futureList.remove(future);  
  44.                     numThread--;  
  45.                     System.out.println(result);  
  46.                     //此处必须break,否则会抛出并发修改异常。(也可以通过将futureList声明为CopyOnWriteArrayList类型解决)  
  47.                     break;  
  48.                 }  
  49.             }  
  50.         }  
  51.     }  
  52. }  

 方式二:

第一种方式显得比较繁琐,通过使用ExecutorCompletionService,则可以达到代码最简化的效果。

Java代码  收藏代码
  1. public class CompletionServiceTest {  
  2.   
  3.     static class Task implements Callable<String>{  
  4.         private int i;  
  5.           
  6.         public Task(int i){  
  7.             this.i = i;  
  8.         }  
  9.   
  10.         @Override  
  11.         public String call() throws Exception {  
  12.             Thread.sleep(10000);  
  13.             return Thread.currentThread().getName() + "执行完任务:" + i;  
  14.         }     
  15.     }  
  16.       
  17.     public static void main(String[] args) throws InterruptedException, ExecutionException{  
  18.         testExecutorCompletionService();  
  19.     }  
  20.       
  21.     private static void testExecutorCompletionService() throws InterruptedException, ExecutionException{  
  22.         int numThread = 5;  
  23.         ExecutorService executor = Executors.newFixedThreadPool(numThread);  
  24.         CompletionService<String> completionService = new ExecutorCompletionService<String>(executor);  
  25.         for(int i = 0;i<numThread;i++ ){  
  26.             completionService.submit(new CompletionServiceTest.Task(i));  
  27.         }  
  28. }  
  29.           
  30.         for(int i = 0;i<numThread;i++ ){       
  31.             System.out.println(completionService.take().get());  
  32.         }  
  33.           
  34.     }  

 

ExecutorCompletionService分析:

 CompletionService是Executor和BlockingQueue的结合体。

Java代码  收藏代码
  1. public ExecutorCompletionService(Executor executor) {  
  2.         if (executor == null)  
  3.             throw new NullPointerException();  
  4.         this.executor = executor;  
  5.         this.aes = (executor instanceof AbstractExecutorService) ?  
  6.             (AbstractExecutorService) executor : null;  
  7.         this.completionQueue = new LinkedBlockingQueue<Future<V>>();  
  8.     }  

 任务的提交和执行都是委托给Executor来完成。当提交某个任务时,该任务首先将被包装为一个QueueingFuture,

Java代码  收藏代码
  1. public Future<V> submit(Callable<V> task) {  
  2.         if (task == nullthrow new NullPointerException();  
  3.         RunnableFuture<V> f = newTaskFor(task);  
  4.         executor.execute(new QueueingFuture(f));  
  5.         return f;  
  6.     }  

 QueueingFutureFutureTask的一个子类,通过改写该子类的done方法,可以实现当任务完成时,将结果放入到BlockingQueue中。

 

Java代码  收藏代码
  1. private class QueueingFuture extends FutureTask<Void> {  
  2.         QueueingFuture(RunnableFuture<V> task) {  
  3.             super(task, null);  
  4.             this.task = task;  
  5.         }  
  6.         protected void done() { completionQueue.add(task); }  
  7.         private final Future<V> task;  
  8.     }  

 而通过使用BlockingQueue的take或poll方法,则可以得到结果。在BlockingQueue不存在元素时,这两个操作会阻塞,一旦有结果加入,则立即返回。

Java代码  收藏代码
  1. public Future<V> take() throws InterruptedException {  
  2.     return completionQueue.take();  
  3. }  
  4.   
  5. public Future<V> poll() {  
  6.     return completionQueue.poll();  
  7. }  
分享到:
评论

相关推荐

    32 请按到场顺序发言—Completion Service详解.pdf

    在使用`ExecutorCompletionService`时,我们需要创建一个`ExecutorService`实例,然后将这个`ExecutorService`传递给`ExecutorCompletionService`的构造函数。接着,我们可以向`ExecutorCompletionService`提交任务...

    java多线程编程实践

    `ExecutorCompletionService`结合了`ExecutorService`和`BlockingQueue`的功能,主要用于管理和监控异步任务的执行结果。 #### 三、锁机制 在多线程编程中,锁是确保数据完整性和一致性的重要手段。`java.util....

    JAVA课程学习笔记.doc

    - `java.util.concurrent.CompletionService`:允许获取执行任务的结果,`ExecutorCompletionService` 是其具体实现,结合了 `ExecutorService` 和 `Future`。 5. 线程池执行原理 线程池的执行过程主要包括任务提交...

    ExecutorService与CompletionService对比详解.docx

    在示例中,创建了一个ExecutorCompletionService实例,它继承自CompletionService并且使用ExecutorService作为底层的执行器。提交任务的方式与ExecutorService类似,但获取结果时,不再直接从列表中获取Future,而是...

    java多线程并发

    `ExecutorCompletionService`是`java.util.concurrent`包提供的一个类,它结合了`ExecutorService`和`BlockingQueue`的功能,用于管理和获取已完成的任务结果。 综上所述,Java中的多线程并发机制非常强大,不仅...

    通过多线程任务处理大批量耗时业务并返回结果

    4. `CompletionService`:可能是`ExecutorCompletionService`,它结合了`ExecutorService`和`BlockingQueue`的功能。我们可以使用`CompletionService.take()`方法获取下一个已完成的任务的结果,而不必等待所有任务...

    JAVA并发编程实践-线程的关闭与取消-学习笔记

    10. **TrackingExecutor任务跟踪**:为了确保任务的正常结束,可以使用`ExecutorCompletionService`来跟踪任务的完成情况,并在必要时取消未完成的任务。 11. **处理异常的线程终止**:线程异常终止时,需要正确...

    Java线程池学习资料-全

    `ThreadPoolExecutor`的`submit()`返回`Future`对象,而`ExecutorCompletionService`的`submit()`除了返回`Future`,还支持批量处理结果。 当线程池中的线程抛出异常时,如果使用`submit()`,异常会被捕获并封装在`...

    Java——JUC

    7. **ExecutorCompletionService** - 一个基于`ExecutorService`的增强版服务,用于管理一组异步任务的执行和结果收集。 8. **ScheduledExecutorService** - 支持定时及周期性任务执行的接口,如`...

    springboot集成amazon aws s3对象存储sdk(javav2)

    CompletionService&lt;PartETag&gt; completionService = new ExecutorCompletionService(executor); for (int i = 1; i ; i++) { UploadPartResponse response = s3Client.uploadPart(uploadRequestBuilder.part...

    ThreadDemo.rar

    - `ExecutorCompletionService`:用于管理一组异步任务,等待任务完成并获取结果。 - `ForkJoinPool`和`RecursiveTask`/`RecursiveAction`:基于工作窃取算法的并行计算框架。 8. **线程中断和守护线程**: 使用...

    jvaa面试宝典

    - 并发工具类:Semaphore、CyclicBarrier、CountDownLatch、ExecutorCompletionService等。 - Future和Callable接口:理解异步计算,以及如何获取结果。 通过深入学习这些知识点,Java开发者可以更好地准备面试,...

    java.util.concurrent.uml.pdf

    ExecutorCompletionService类是其实现,它利用线程池执行任务,并帮助开发者获取已经完成的任务结果。 Runnable和Callable是两种任务类型。Runnable是任务的一个简单的执行对象,没有返回值。而Callable接口类似于...

    Callable,Future的使用方式

    Callable,Future的使用方式,里面使用了三种使用方式分别是FutureTask,ExecutorService,ExecutorCompletionService

    多线程相关代码(新)

    包括阻塞队列、阻塞栈、ExecutorService、Future、ExecutorCompletionService、死锁、join、重入锁、读写锁、多线程抢票、信号量、signal/await、ThreadLocal等的实例。

    浅谈Java多线程处理中Future的妙用(附源码)

    CompletionService&lt;Object&gt; ecs = new ExecutorCompletionService(executor); for (int i = 0; i ; i++) { final Integer t = data[i]; ecs.submit(new Callable() { public Object call() { try { Thread....

    java并发框架源码-notes:记录各种学习笔记(Java、算法、框架、数据库、并发、源码...)

    `ExecutorCompletionService`用于批量处理完成的任务,提高效率。 8. **框架源码分析**: 分析如`Akka`、`Quasar`或`Disruptor`等并发框架的源码,可以深入理解如何在Java中构建高效的并发系统,学习其设计思想和...

    java抢票系统源码-POS:邮政

    5. **异步编程**:Java 8引入了CompletableFuture和ExecutorCompletionService等工具,使得开发者能更高效地处理异步任务,提高系统性能。 6. **Web框架**:为了简化开发,项目可能使用Spring Boot或Struts等Web...

    CoreJava:核心Java

    14. **并发编程**:深入研究并发工具类(如CountDownLatch, CyclicBarrier, Semaphore, ExecutorCompletionService等),以及并发容器(如ConcurrentHashMap, CopyOnWriteArrayList等)。 15. **垃圾回收与内存管理...

Global site tag (gtag.js) - Google Analytics