最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java多線程 CompletionService

 更新時間:2021年10月27日 17:26:28   作者:冬日毛毛雨  
這篇文章主要介紹了Java多線程 CompletionService,CompletionService用于提交一組Callable任務,其take方法返回已完成的一個Callable任務對應的Future對象,需要的朋友可以參考一下文章詳細內容

1 CompletionService介紹

CompletionService用于提交一組Callable任務,其take方法返回已完成的一個Callable任務對應的Future對象。
如果你向Executor提交了一個批處理任務,并且希望在它們完成后獲得結果。為此你可以將每個任務的Future保存進一個集合,然后循環(huán)這個集合調用Futureget()取出數據。幸運的是CompletionService幫你做了這件事情。
CompletionService整合了ExecutorBlockingQueue的功能。你可以將Callable任務提交給它去執(zhí)行,然后使用類似于隊列中的take和poll方法,在結果完整可用時獲得這個結果,像一個打包的Future
CompletionService的take返回的future是哪個先完成就先返回哪一個,而不是根據提交順序。

2 CompletionService源碼分析

首先看一下 構造方法:

   public ExecutorCompletionService(Executor executor) {
        if (executor == null)
            throw new NullPointerException();
        this.executor = executor;
        this.aes = (executor instanceof AbstractExecutorService) ?
            (AbstractExecutorService) executor : null;
        this.completionQueue = new LinkedBlockingQueue<Future<V>>();
    }

構造法方法主要初始化了一個阻塞隊列,用來存儲已完成的task任務。

然后看一下 completionService.submit 方法:

    public Future<V> submit(Callable<V> task) {
        if (task == null) throw new NullPointerException();
        RunnableFuture<V> f = newTaskFor(task);
        executor.execute(new QueueingFuture(f));
        return f;
    }

    public Future<V> submit(Runnable task, V result) {
        if (task == null) throw new NullPointerException();
        RunnableFuture<V> f = newTaskFor(task, result);
        executor.execute(new QueueingFuture(f));
        return f;
    }

可以看到,callable任務被包裝成QueueingFuture,而 QueueingFutureFutureTask的子類,所以最終執(zhí)行了FutureTask中的run()方法。

來看一下該方法:

 public void run() {
 //判斷執(zhí)行狀態(tài),保證callable任務只被運行一次
    if (state != NEW ||
        !UNSAFE.compareAndSwapObject(this, runnerOffset,
                                     null, Thread.currentThread()))
        return;
    try {
        Callable<V> c = callable;
        if (c != null && state == NEW) {
            V result;
            boolean ran;
            try {
            //這里回調我們創(chuàng)建的callable對象中的call方法
                result = c.call();
                ran = true;
            } catch (Throwable ex) {
                result = null;
                ran = false;
                setException(ex);
            }
            if (ran)
            //處理執(zhí)行結果
                set(result);
        }
    } finally {
        runner = null;
        // state must be re-read after nulling runner to prevent
        // leaked interrupts
        int s = state;
        if (s >= INTERRUPTING)
            handlePossibleCancellationInterrupt(s);
    }
}

可以看到在該 FutureTask 中執(zhí)行run方法,最終回調自定義的callable中的call方法,執(zhí)行結束之后,

通過 set(result) 處理執(zhí)行結果:

    /**
     * Sets the result of this future to the given value unless
     * this future has already been set or has been cancelled.
     *
     * <p>This method is invoked internally by the {@link #run} method
     * upon successful completion of the computation.
     *
     * @param v the value
     */
    protected void set(V v) {
        if (UNSAFE.compareAndSwapInt(this, stateOffset, NEW, COMPLETING)) {
            outcome = v;
            UNSAFE.putOrderedInt(this, stateOffset, NORMAL); // final state
            finishCompletion();
        }
    }

繼續(xù)跟進finishCompletion()方法,在該方法中找到 done()方法:

protected void done() { completionQueue.add(task); }

可以看到該方法只做了一件事情,就是將執(zhí)行結束的task添加到了隊列中,只要隊列中有元素,我們調用take()方法時就可以獲得執(zhí)行的結果。
到這里就已經清晰了,異步非阻塞獲取執(zhí)行結果的實現(xiàn)原理其實就是通過隊列來實現(xiàn)的,FutureTask將執(zhí)行結果放到隊列中,先進先出,線程執(zhí)行結束的順序就是獲取結果的順序。

CompletionService實際上可以看做是ExecutorBlockingQueue的結合體。CompletionService在接收到要執(zhí)行的任務時,通過類似BlockingQueue的put和take獲得任務執(zhí)行的結果。CompletionService的一個實現(xiàn)是ExecutorCompletionService,ExecutorCompletionService把具體的計算任務交給Executor完成。

在實現(xiàn)上,ExecutorCompletionService在構造函數中會創(chuàng)建一個BlockingQueue(使用的基于鏈表的無界隊列LinkedBlockingQueue),該BlockingQueue的作用是保存Executor執(zhí)行的結果。當計算完成時,調用FutureTask的done方法。當提交一個任務到ExecutorCompletionService時,首先將任務包裝成QueueingFuture,它是FutureTask的一個子類,然后改寫FutureTask的done方法,之后把Executor執(zhí)行的計算結果放入BlockingQueue中。

QueueingFuture的源碼如下:

    /**
     * FutureTask extension to enqueue upon completion
     */
    private class QueueingFuture extends FutureTask<Void> {
        QueueingFuture(RunnableFuture<V> task) {
            super(task, null);
            this.task = task;
        }
        protected void done() { completionQueue.add(task); }
        private final Future<V> task;
    }

3 CompletionService實現(xiàn)任務

public class CompletionServiceTest {
    public static void main(String[] args) {

        ExecutorService threadPool = Executors.newFixedThreadPool(10);
        CompletionService<Integer> completionService = new ExecutorCompletionService<Integer>(threadPool);
        for (int i = 1; i <=10; i++) {
            final int seq = i;
            completionService.submit(new Callable<Integer>() {
                @Override
                public Integer call() throws Exception {

                    Thread.sleep(new Random().nextInt(5000));

                    return seq;
                }
            });
        }
        threadPool.shutdown();
        for (int i = 0; i < 10; i++) {
            try {
                System.out.println(
                        completionService.take().get());
            } catch (InterruptedException e) {
                e.printStackTrace();
            } catch (ExecutionException e) {
                e.printStackTrace();
            }
        }

    }
}

7
3
9
8
1
2
4
6
5
10

4 CompletionService總結

相比ExecutorService,CompletionService可以更精確和簡便地完成異步任務的執(zhí)行
CompletionService的一個實現(xiàn)是ExecutorCompletionService,它是ExecutorBlockingQueue功能的融合體,Executor完成計算任務,BlockingQueue負責保存異步任務的執(zhí)行結果
在執(zhí)行大量相互獨立和同構的任務時,可以使用CompletionService
CompletionService可以為任務的執(zhí)行設置時限,主要是通過BlockingQueuepoll(long time,TimeUnit unit)為任務執(zhí)行結果的取得限制時間,如果沒有完成就取消任務

到此這篇關于Java多線程 CompletionService的文章就介紹到這了,更多相關Java多線程CompletionService內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • java異常繼承何類,運行時異常與一般異常的區(qū)別(詳解)

    java異常繼承何類,運行時異常與一般異常的區(qū)別(詳解)

    下面小編就為大家?guī)硪黄猨ava異常繼承何類,運行時異常與一般異常的區(qū)別(詳解)。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-11-11
  • 基于Intellij Idea亂碼的解決方法

    基于Intellij Idea亂碼的解決方法

    下面小編就為大家分享一篇基于Intellij Idea亂碼的解決方法,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2018-03-03
  • Java 并發(fā)編程創(chuàng)建線程的四種方式示例詳解

    Java 并發(fā)編程創(chuàng)建線程的四種方式示例詳解

    本文對比Java創(chuàng)建線程的四種方式:繼承Thread類(單繼承限制)、實現(xiàn)Runnable接口(靈活推薦)、實現(xiàn)Callable接口(支持返回值)、使用線程池(最佳實踐,資源管理高效),推薦優(yōu)先使用線程池提升系統(tǒng)穩(wěn)定性與資源利用率,感興趣的朋友一起看看吧
    2025-09-09
  • Spring中的Bean對象初始化詳解

    Spring中的Bean對象初始化詳解

    這篇文章主要介紹了Spring中的Bean對象初始化詳解,我們也可以根據Java語言的特性猜測到其很有可能是通過反射機制來完成Bean的初始化操作,接下來我們一步一步的剖析Spring對Bean的初始化操作,需要的朋友可以參考下
    2023-12-12
  • java多線程CountDownLatch與線程池ThreadPoolExecutor/ExecutorService案例

    java多線程CountDownLatch與線程池ThreadPoolExecutor/ExecutorService案

    這篇文章主要介紹了java多線程CountDownLatch與線程池ThreadPoolExecutor/ExecutorService案例,
    2021-02-02
  • 使用nacos增加修改配置實時生效方式

    使用nacos增加修改配置實時生效方式

    文章介紹了如何在Nacos配置中心實現(xiàn)動態(tài)配置,包括添加或修改配置、使用`@RefreshScope`注解刷新Bean屬性、擴展動態(tài)配置參數等
    2026-01-01
  • SpringCache快速使用及入門案例

    SpringCache快速使用及入門案例

    Spring Cache 是Spring 提供的一整套的緩存解決方案,它不是具體的緩存實現(xiàn),本文主要介紹了SpringCache快速使用及入門案例,感興趣的可以了解一下
    2023-08-08
  • 基于SpringBoot打造一個通用CLI命令系統(tǒng)

    基于SpringBoot打造一個通用CLI命令系統(tǒng)

    在日常開發(fā)中,某些情況下可能需要為服務提供一個命令行工具(CLI),方便運維、調試或者遠程調用業(yè)務接口,下面我們就來使用SpringBoot打造一個通用的CLI命令系統(tǒng)吧
    2025-12-12
  • SpringBoot整合Kafka完成生產消費的方案

    SpringBoot整合Kafka完成生產消費的方案

    網上找了很多管理kafka整合springboot的教程,但是很多都沒辦法應用到生產環(huán)境,很多配置都是缺少,或者不正確的,只能當個demo,所以本文給大家介紹了SpringBoot整合Kafka完成生產消費的方案,需要的朋友可以參考下
    2024-12-12
  • Java線程數究竟設多少合理

    Java線程數究竟設多少合理

    這篇文章主要介紹了Java線程數究竟設多少合理,對線程感興趣的同學,可以參考下
    2021-04-04

最新評論

永新县| 乳山市| 郯城县| 南安市| 德格县| 台北市| 宁乡县| 郯城县| 比如县| 贞丰县| 天台县| 乌苏市| 从江县| 德庆县| 潼关县| 江川县| 天长市| 房山区| 上犹县| 囊谦县| 监利县| 浮梁县| 嘉义市| 布拖县| 平和县| 鹤山市| 丰镇市| 赤峰市| 固安县| 东辽县| 新和县| 库车县| 彭州市| 渝中区| 南城县| 三江| 方城县| 阜康市| 镇安县| 文昌市| 邵武市|