Java實現(xiàn)CompletionService并發(fā)編排消費任務(wù)
假如出現(xiàn)了這種情況:RocketMQ 批量拉取了 128 條消息,但消費端是一條一條串行處理的。128 條消息,每條 50 毫秒,一輪消費就要 6.4 秒。
可以直接用批量消費嗎?
串行消費的瓶頸
RocketMQ 的 setConsumeMessageBatchMaxSize(128) 讓你一次拉 128 條消息過來,但如果你還是逐條同步處理,那批量拉取的意義就只剩「減少網(wǎng)絡(luò)往返」了。處理速度的瓶頸,一條都沒解開。
直覺上的解法很簡單,交給線程池并發(fā)執(zhí)行嘛。但這里藏著一個很容易忽略的問題。
如果異步線程執(zhí)行失敗了,RocketMQ 的 Broker 是不知道的。主線程已經(jīng)返回了 CONSUME_SUCCESS,Broker 提交了偏移量,那條失敗的消息就丟了。或者更常見的情況,主線程根本不知道哪些子線程成功了、哪些失敗了,只能盲目地返回成功或重試。
我們需要一個辦法,在主線程中感知每一個子線程的執(zhí)行結(jié)果。全部成功才返回 CONSUME_SUCCESS,有一條失敗就返回 RECONSUME_LATER 讓 RocketMQ 整體重發(fā)。
這不是「能不能并發(fā)」的問題,是「并發(fā)了能不能控」的問題。
CompletionService,先完成先取結(jié)果的編排器
Java 并發(fā)包里有一個接口叫 CompletionService,位于 java.util.concurrent,專門干這件事,批量提交異步任務(wù),按完成順序逐個取結(jié)果。
它的唯一實現(xiàn)類是 ExecutorCompletionService,用法就三步。
- 提交所有任務(wù)到線程池。
- 循環(huán)取結(jié)果,發(fā)現(xiàn)失敗立即標記。
- 全部成功返回
CONSUME_SUCCESS,否則返回RECONSUME_LATER。
代碼邏輯如下。
@Override
public void prepareStart(DefaultMQPushConsumer consumer) {
consumer.setPullInterval(1000);
consumer.setConsumeMessageBatchMaxSize(128);
consumer.setPullBatchSize(64);
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
log.info("NewBuyBatchMsgListener receive message size: {}", msgs.size());
CompletionService<Boolean> completionService = new ExecutorCompletionService<>(executor);
List<Future<Boolean>> futures = new ArrayList<>();
// 1. 提交所有任務(wù)
msgs.forEach(messageExt -> {
Callable<Boolean> task = () -> {
try {
OrderCreateRequest orderCreateRequest = JSON.parseObject(JSON.parseObject(messageExt.getBody()).getString("body"), OrderCreateRequest.class);
return doNewBuyExecute(orderCreateRequest);
} catch (Exception e) {
log.error("Task failed", e);
return false; // 標記失敗
}
};
futures.add(completionService.submit(task));
});
// 2. 檢查結(jié)果
boolean allSuccess = true;
try {
for (int i = 0; i < msgs.size(); i++) {
Future<Boolean> future = completionService.take();
if (!future.get()) { // 3. 發(fā)現(xiàn)一個失敗立即終止
allSuccess = false;
break;
}
}
} catch (Exception e) {
allSuccess = false;
}
// 3. 根據(jù)結(jié)果返回消費狀態(tài)
return allSuccess ? ConsumeConcurrentlyStatus.CONSUME_SUCCESS
: ConsumeConcurrentlyStatus.RECONSUME_LATER;
});
}
128 條消息并發(fā)處理,總耗時從 6.4 秒降到最慢那條的耗時,通常幾十毫秒就夠了。
那么,「CompletionService 底層是怎么做到先完成先取的?」
底層原理,BlockingQueue + QueueingFuture
ExecutorCompletionService 的源碼非常精煉,核心就三個成員變量。
private final Executor executor; private final AbstractExecutorService aes; private final BlockingQueue<Future<V>> completionQueue;
executor 是你傳入的線程池,completionQueue 是一個 LinkedBlockingQueue,用來存放已完成任務(wù)的 Future 對象。
關(guān)鍵在 submit 方法里。當你調(diào)用 completionService.submit(task) 時,它并沒有直接把 task 丟給線程池,而是先包裝了一層。
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;
}
QueueingFuture 繼承自 FutureTask,重寫了 done() 方法。done() 是 FutureTask 提供的鉤子,任務(wù)無論正常完成還是異常終止,都會回調(diào)這個方法。
所以整個流程是這樣的。你 submit 一個任務(wù),它被包裝成 QueueingFuture 交給線程池執(zhí)行。任務(wù)跑完的那一刻,done() 觸發(fā),把對應的 Future 塞進 completionQueue。你調(diào) take(),就是從 completionQueue 里阻塞地取一個出來。
誰先完成,誰的 Future 先入隊,你就先取到誰的結(jié)果。提交順序和完成順序解耦了。
我覺得這個設(shè)計很漂亮。它沒有用任何鎖排序、沒有用優(yōu)先隊列、沒有用回調(diào)鏈,就是最樸素的「生產(chǎn)者往隊列里放,消費者從隊列里取」。BlockingQueue 天然線程安全,生產(chǎn)消費解耦,簡單到幾乎不可能出錯。
那 CompletableFuture 呢?
Java 8 引入了 CompletableFuture,同樣是處理異步任務(wù)的利器。它和 CompletionService 解決的問題有重疊,但設(shè)計哲學完全不同。
| 維度 | CompletionService | CompletableFuture |
|---|---|---|
| 引入版本 | Java 5 | Java 8 |
| 核心機制 | BlockingQueue,先完成先取 | 回調(diào)鏈,任務(wù)間可編排依賴 |
| 結(jié)果獲取 | take() 阻塞等待下一個完成 | thenApply() / thenCompose() 非阻塞回調(diào) |
| 任務(wù)關(guān)系 | 批量獨立任務(wù),互不依賴 | 可描述 A 完成后執(zhí)行 B、A 和 B 都完成后執(zhí)行 C |
| 異常處理 | 在 Future.get() 時拋 ExecutionException | exceptionally() / handle() 流式處理 |
| 適用場景 | 批量同構(gòu)任務(wù),只關(guān)心結(jié)果是否全部成功 | 異步流程編排,任務(wù)間有依賴和組合關(guān)系 |
一句話總結(jié),CompletionService 是「批量并發(fā),按完成順序收結(jié)果」,CompletableFuture 是「異步編排,按依賴關(guān)系串流程」。
回到我們這個 RocketMQ 批量消費的場景,128 條消息之間沒有任何依賴關(guān)系,我們只關(guān)心「全部成功還是有一個失敗」。這就是 CompletionService 的主場。
如果你用 CompletableFuture 來寫,也能做,但你需要自己維護一個 CompletableFuture.allOf() 來等全部完成,然后再遍歷檢查結(jié)果。代碼更啰嗦,而且 allOf 會等所有任務(wù)都完成才能繼續(xù),哪怕第 2 條消息就失敗了,你也要等剩下 126 條跑完才能返回。CompletionService 的 take() 則是逐個檢查,發(fā)現(xiàn)失敗立即 break,省下了不必要的等待。
當然,如果你的場景是「查商品信息,再根據(jù)商品查庫存和價格,最后組裝結(jié)果」,任務(wù)之間有明確的先后依賴,那 CompletableFuture 的鏈式編排就比 CompletionService 的隊列取值優(yōu)雅得多。
工具沒有好壞,只有合不合適。
并發(fā)不是目的,可控才是
串行消費慢,直覺反應是加并發(fā)。但加了并發(fā)之后,如果主線程無法感知子線程的成敗,那并發(fā)就不是加速,是埋雷。消息丟了都不知道。
CompletionService 解決的不是「怎么并發(fā)」的問題,而是「并發(fā)了怎么收場」的問題。它用最樸素的 BlockingQueue 機制,讓主線程能按完成順序逐個檢查結(jié)果,發(fā)現(xiàn)異常立即止損。
我覺得并發(fā)編程最難的部分從來不是「怎么讓任務(wù)跑起來」,而是「跑起來之后怎么確保結(jié)果可控」。CompletionService 給了一個很干凈的答案。
到此這篇關(guān)于Java實現(xiàn)CompletionService并發(fā)編排消費任務(wù)的文章就介紹到這了,更多相關(guān)Java CompletionService并發(fā)編排內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Spring Aop 如何獲取參數(shù)名參數(shù)值
這篇文章主要介紹了Spring Aop 如何獲取參數(shù)名參數(shù)值的操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-07-07
java多線程編程之使用runnable接口創(chuàng)建線程
實現(xiàn)Runnable接口的類必須使用Thread類的實例才能創(chuàng)建線程,通過Runnable接口創(chuàng)建線程分為以下兩步2014-01-01
SpringBoot2 整合Ehcache組件,輕量級緩存管理的原理解析
這篇文章主要介紹了SpringBoot2 整合Ehcache組件,輕量級緩存管理,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2020-08-08
Spring線程池ThreadPoolTaskExecutor配置與實踐方式
文章介紹了Spring框架的ThreadPoolTaskExecutor,其主要功能包括線程池管理、任務(wù)執(zhí)行及高級特性,核心參數(shù)有核心線程數(shù)、最大線程數(shù)、隊列容量等,文中還詳細解釋了拒絕策略、線程上下文類加載器等高級功能,并提供了配置示例和使用建議,幫助開發(fā)者有效管理線程池2026-03-03
學習不同 Java.net 語言中類似的函數(shù)結(jié)構(gòu)
這篇文章主要介紹了學習不同 Java.net 語言中類似的函數(shù)結(jié)構(gòu),函數(shù)式編程語言包含多個系列的常見函數(shù)。但開發(fā)人員有時很難在語言之間進行切換,因為熟悉的函數(shù)具有不熟悉的名稱。函數(shù)式語言傾向于基于函數(shù)范例來命名這些常見函數(shù)。,需要的朋友可以參考下2019-06-06

