摘要:這個方法返回與等待所有返回等待多個返回取多個當中最快的一個返回等待多個當中最快的一個返回二詳解終極指南并發(fā)編程中的風格
thenApply(等待并轉(zhuǎn)化future)
@Test public void testThen() throws ExecutionException, InterruptedException { CompletableFuturethenAccept與thenRun(監(jiān)聽future完成)f1 = CompletableFuture.supplyAsync(() -> { return "zero"; }, executor); CompletableFuture f2 = f1.thenApply(new Function () { @Override public Integer apply(String t) { System.out.println(2); return Integer.valueOf(t.length()); } }); CompletableFuture f3 = f2.thenApply(r -> r * 2.0); System.out.println(f3.get()); }
/** * future完成處理,可獲取結(jié)果 */ @Test public void testThenAccept(){ CompletableFuturethenCompose(flatMap future)f1 = CompletableFuture.supplyAsync(() -> { return "zero"; }, executor); f1.thenAccept(e -> { System.out.println("get result:"+e); }); } /** * future完成處理 */ @Test public void testThenRun(){ CompletableFuture f1 = CompletableFuture.supplyAsync(() -> { return "zero"; }, executor); f1.thenRun(new Runnable() { @Override public void run() { System.out.println("finished"); } }); }
/** * compose相當于flatMap,避免CompletableFuturethenCombine與thenAcceptBoth>這種 * @throws ExecutionException * @throws InterruptedException */ @Test public void testThenCompose() throws ExecutionException, InterruptedException { ExecutorService executor = Executors.newFixedThreadPool(5); CompletableFuture f1 = CompletableFuture.supplyAsync(() -> { return "zero"; }, executor); CompletableFuture > f4 = f1.thenApply(CompletableFutureTest::calculate); System.out.println("f4.get:"+f4.get().get()); CompletableFuture f5 = f1.thenCompose(CompletableFutureTest::calculate); System.out.println("f5.get:"+f5.get()); System.out.println(f1.get()); } public static CompletableFuture calculate(String input) { ExecutorService executor = Executors.newFixedThreadPool(5); CompletableFuture future = CompletableFuture.supplyAsync(() -> { System.out.println(input); return input + "---" + input.length(); }, executor); return future; }
thenCombine(組合兩個future,有返回值)
/** * thenCombine用于組合兩個并發(fā)的任務(wù),產(chǎn)生新的future有返回值 * @throws ExecutionException * @throws InterruptedException */ @Test public void testThenCombine() throws ExecutionException, InterruptedException { CompletableFuturef1 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f1 start to sleep at:"+System.currentTimeMillis()); Thread.sleep(1000); System.out.println("f1 finish sleep at:"+System.currentTimeMillis()); } catch (InterruptedException e) { e.printStackTrace(); } return "zero"; }, executor); CompletableFuture f2 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f2 start to sleep at:"+System.currentTimeMillis()); Thread.sleep(3000); System.out.println("f2 finish sleep at:"+System.currentTimeMillis()); } catch (InterruptedException e) { e.printStackTrace(); } return "hello"; }, executor); CompletableFuture reslutFuture = f1.thenCombine(f2, new BiFunction () { @Override public String apply(String t, String u) { System.out.println("f3 start to combine at:"+System.currentTimeMillis()); return t.concat(u); } }); System.out.println(reslutFuture.get());//zerohello System.out.println("finish combine at:"+System.currentTimeMillis()); }
thenAcceptBoth(組合兩個future,沒有返回值)
/** * thenAcceptBoth用于組合兩個并發(fā)的任務(wù),產(chǎn)生新的future沒有返回值 * @throws ExecutionException * @throws InterruptedException */ @Test public void testThenAcceptBoth() throws ExecutionException, InterruptedException { CompletableFutureapplyToEither與acceptEitherf1 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f1 start to sleep at:"+System.currentTimeMillis()); TimeUnit.SECONDS.sleep(1); System.out.println("f1 stop sleep at:"+System.currentTimeMillis()); } catch (Exception e) { // TODO Auto-generated catch block e.printStackTrace(); } return "zero"; }, executor); CompletableFuture f2 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f2 start to sleep at:"+System.currentTimeMillis()); TimeUnit.SECONDS.sleep(3); System.out.println("f2 stop sleep at:"+System.currentTimeMillis()); } catch (Exception e) { // TODO Auto-generated catch block e.printStackTrace(); } return "hello"; }, executor); CompletableFuture reslutFuture = f1.thenAcceptBoth(f2, new BiConsumer () { @Override public void accept(String t, String u) { System.out.println("f3 start to accept at:"+System.currentTimeMillis()); System.out.println(t + " over"); System.out.println(u + " over"); } }); System.out.println(reslutFuture.get()); System.out.println("finish accept at:"+System.currentTimeMillis()); }
applyToEither(取2個future中最先返回的,有返回值)
/** * 當任意一個CompletionStage 完成的時候,fn 會被執(zhí)行,它的返回值會當做新的CompletableFuture的計算結(jié)果 * @throws ExecutionException * @throws InterruptedException */ @Test public void testApplyToEither() throws ExecutionException, InterruptedException { CompletableFuturef1 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f1 start to sleep at:"+System.currentTimeMillis()); TimeUnit.SECONDS.sleep(5); System.out.println("f1 stop sleep at:"+System.currentTimeMillis()); } catch (Exception e) { // TODO Auto-generated catch block e.printStackTrace(); } return "fromF1"; }, executor); CompletableFuture f2 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f2 start to sleep at:"+System.currentTimeMillis()); TimeUnit.SECONDS.sleep(2); System.out.println("f2 stop sleep at:"+System.currentTimeMillis()); } catch (Exception e) { // TODO Auto-generated catch block e.printStackTrace(); } return "fromF2"; }, executor); CompletableFuture reslutFuture = f1.applyToEither(f2,i -> i.toString()); System.out.println(reslutFuture.get()); //should not be null , wait for complete }
acceptEither(取2個future中最先返回的,無返回值)
/** * 取其中返回最快的一個 * 當任意一個CompletionStage 完成的時候,action 這個消費者就會被執(zhí)行。這個方法返回 CompletableFutureallOf與anyOf*/ @Test public void testAcceptEither() throws ExecutionException, InterruptedException { CompletableFuture f1 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f1 start to sleep at:"+System.currentTimeMillis()); TimeUnit.SECONDS.sleep(3); System.out.println("f1 stop sleep at:"+System.currentTimeMillis()); } catch (Exception e) { // TODO Auto-generated catch block e.printStackTrace(); } return "zero"; }, executor); CompletableFuture f2 = CompletableFuture.supplyAsync(() -> { try { System.out.println("f2 start to sleep at:"+System.currentTimeMillis()); TimeUnit.SECONDS.sleep(5); System.out.println("f2 stop sleep at:"+System.currentTimeMillis()); } catch (Exception e) { // TODO Auto-generated catch block e.printStackTrace(); } return "hello"; }, executor); CompletableFuture reslutFuture = f1.acceptEither(f2,r -> { System.out.println("quicker result:"+r); }); reslutFuture.get(); //should be null , wait for complete }
allOf(等待所有future返回)
/** * 等待多個future返回 */ @Test public void testAllOf() throws InterruptedException { List> futures = IntStream.range(1,10) .mapToObj(i -> longCost(i)).collect(Collectors.toList()); final CompletableFuture allCompleted = CompletableFuture.allOf(futures.toArray(new CompletableFuture[]{})); allCompleted.thenRun(() -> { futures.stream().forEach(future -> { try { System.out.println("get future at:"+System.currentTimeMillis()+", result:"+future.get()); } catch (InterruptedException | ExecutionException e) { e.printStackTrace(); } }); }); Thread.sleep(100000); //wait }
anyOf(取多個future當中最快的一個返回)
/** * 等待多個future當中最快的一個返回 * @throws InterruptedException */ @Test public void testAnyOf() throws InterruptedException { Listdoc> futures = IntStream.range(1,10) .mapToObj(i -> longCost(i)).collect(Collectors.toList()); final CompletableFuture
CompletableFuture(二)
Java CompletableFuture 詳解
Java 8:CompletableFuture終極指南
Java 8: Definitive guide to CompletableFuture
并發(fā)編程 | JDK 1.8中的CompletableFuture | FRP風格
文章版權(quán)歸作者所有,未經(jīng)允許請勿轉(zhuǎn)載,若此文章存在違規(guī)行為,您可以聯(lián)系管理員刪除。
轉(zhuǎn)載請注明本文地址:http://systransis.cn/yun/66708.html
摘要:方法接收的是的實例,但是它沒有返回值方法是函數(shù)式接口,無參數(shù),會返回一個結(jié)果這兩個方法是的升級,表示讓任務(wù)在指定的線程池中執(zhí)行,不指定的話,通常任務(wù)是在線程池中執(zhí)行的。該的接口是在線程使用舊的接口,它不允許返回值。 簡介 作為Java 8 Concurrency API改進而引入,本文是CompletableFuture類的功能和用例的介紹。同時在Java 9 也有對Completab...
摘要:組合式異步編程最近這些年,兩種趨勢不斷地推動我們反思我們設(shè)計軟件的方式。第章中介紹的分支合并框架以及并行流是實現(xiàn)并行處理的寶貴工具它們將一個操作切分為多個子操作,在多個不同的核甚至是機器上并行地執(zhí)行這些子操作。 CompletableFuture:組合式異步編程 最近這些年,兩種趨勢不斷地推動我們反思我們設(shè)計軟件的方式。第一種趨勢和應(yīng)用運行的硬件平臺相關(guān),第二種趨勢與應(yīng)用程序的架構(gòu)相關(guān)...
摘要:首先想到的是開啟一個新的線程去做某項工作。再進一步,為了讓新線程可以返回一個值,告訴主線程事情做完了,于是乎粉墨登場。然而提供的方式是主線程主動問詢新線程,要是有個回調(diào)函數(shù)就爽了。極大的提高效率。 showImg(https://segmentfault.com/img/bVbvgBJ?w=1920&h=1200); 引子 為了讓程序更加高效,讓CPU最大效率的工作,我們會采用異步編程...
摘要:內(nèi)部類,用于對和異常進行包裝,從而保證對進行只有一次成功。是取消異常,轉(zhuǎn)換后拋出。判斷是否使用的線程池,在中持有該線程池的引用。 前言 近期作者對響應(yīng)式編程越發(fā)感興趣,在內(nèi)部分享JAVA9-12新特性過程中,有兩處特性讓作者深感興趣:1.JAVA9中的JEP266對并發(fā)編程工具的更新,包含發(fā)布訂閱框架Flow和CompletableFuture加強,其中發(fā)布訂閱框架以java.base...
摘要:項目需求項目中需要優(yōu)化一個接口,這個接口需要拉取個第三方接口,需求延遲時間小于技術(shù)選型是提出的一個支持非阻塞的多功能的,同樣也是實現(xiàn)了接口,是添加的類,用來描述一個異步計算的結(jié)果。對進一步完善,擴展了諸多功能形成了。 項目需求: 項目中需要優(yōu)化一個接口,這個接口需要拉取23個第三方接口,需求延遲時間小于200ms; 技術(shù)選型: CompletableFuture是JDK8提出的一個支持...
閱讀 2276·2021-11-16 11:44
閱讀 650·2019-08-30 15:55
閱讀 3282·2019-08-30 15:52
閱讀 3621·2019-08-30 15:43
閱讀 2205·2019-08-30 11:21
閱讀 445·2019-08-29 12:18
閱讀 1954·2019-08-26 18:15
閱讀 478·2019-08-26 10:32