Java滑動窗口實現(xiàn)線程池并發(fā)度控制的方法詳解
前言
本文深入解析了如何使用 ExecutorCompletionService 和 Guava ListenableFuture 實現(xiàn)并發(fā)度可控的任務執(zhí)行器,幫助你掌握生產級并發(fā)編程的核心技巧。
問題背景
在高并發(fā)場景下,我們可能面臨這樣的需求(僅為例子):
場景:火車票查詢系統(tǒng),用戶輸入北京到上海,需要查詢100個車次的余票信息。
需求:
- 批量查詢:并發(fā)調用100個車次的余票查詢接口
- 限制并發(fā)度:12306限流,單個應用最多10個并發(fā),超過會返回"請稍后再試"
- 盡快返回:不能等全部查完,需要立即返回Future供上層組合(邊查邊展示)
- 按序映射:返回的Future列表要和車次列表一一對應(G1-G100)
傳統(tǒng)方案的痛點:
| 方案 | 問題 |
|---|---|
ExecutorService.invokeAll() | 阻塞等待所有任務完成,無法立即返回 |
CompletableFuture.allOf() | 無法控制并發(fā)度,100個請求同時發(fā)出觸發(fā)限流 |
| 自己用 Semaphore | 代碼復雜,需要手動管理信號量和異常 |
| 分批執(zhí)行(10個一批) | 任務執(zhí)行時間不均勻,短任務等待長任務,資源利用率低 |
優(yōu)化方案:滑動窗口 - 初始提交10個,每完成一個立即補充下一個,保持10個并發(fā)。

代碼實現(xiàn)
/**
* 并發(fā)限制執(zhí)行器
* <p>
* 核心功能:控制任務并發(fā)度,采用滑動窗口策略
* <p>
* 知識點:
* 1. ExecutorCompletionService - 按完成順序獲取任務結果
* 2. ListenableFuture - Guava 可監(jiān)聽 Future
* 3. SettableFuture - 可手動設置結果的 Future
* 4. 滑動窗口并發(fā)控制 - 初始提交N個,每完成一個立即補充下一個
*
* @author 樺說編程
*/
@Value
public class ConcurrentLimitExecutor<V> {
ListeningExecutorService pool;
int parallelism;
BlockingQueue<Future<V>> q;
ExecutorCompletionService<V> cs = new ExecutorCompletionService<>(pool, q);
ListeningExecutorService submitter = MoreExecutors.listeningDecorator(
Executors.newSingleThreadExecutor());
/**
* 提交所有任務,立即返回Future列表
*
* @param tasks 待執(zhí)行的任務列表
* @return Future列表,順序與tasks一致
*/
@SuppressWarnings("unchecked")
public List<ListenableFuture<V>> submitAll(List<Callable<V>> tasks) {
if (tasks.isEmpty()) {
return ImmutableList.of();
}
// 創(chuàng)建結果占位符,數量等于任務數量
List<SettableFuture<V>> result = IntStream.range(0, tasks.size())
.mapToObj(__ -> SettableFuture.<V>create())
.collect(toImmutableList());
// 首批提交數量
int start = Math.min(tasks.size(), parallelism);
// 提交首批任務
for (int i = 0; i < start; i++) {
ListenableFuture<V> f = (ListenableFuture<V>) cs.submit(tasks.get(i));
linkFuture(f, result.get(i));
}
// 異步提交剩余任務
submitter.submit(() -> submitRemaining(tasks, result, start));
return (List<ListenableFuture<V>>) (List<?>) result;
}
/**
* 提交剩余任務(在獨立線程中執(zhí)行)
* <p>
* 流程:
* 1. 等待任意任務完成(cs.take())
* 2. 提交下一個任務
* 3. 將新任務的Future鏈接到result對應位置
* 4. 重復直到所有任務提交完
*/
@SneakyThrows
private void submitRemaining(List<Callable<V>> tasks,
List<SettableFuture<V>> result,
int start) {
int index = start;
int size = tasks.size();
while (index < size) {
// 阻塞等待任意任務完成(釋放一個并發(fā)槽位)
cs.take();
// 提交下一個任務
ListenableFuture<V> f = (ListenableFuture<V>) cs.submit(tasks.get(index));
// 鏈接結果到對應位置
linkFuture(f, result.get(index));
index++;
}
}
private void linkFuture(ListenableFuture<V> from, SettableFuture<V> to) {
to.setFuture(from);
}
}關鍵設計點:
| 設計元素 | 說明 |
|---|---|
| result 數量 | 等于 tasks.size(),每個任務都有對應的 Future |
| linkFuture | 將 cs 返回的 Future 鏈接到 result 中對應位置 |
| 異步補充 | 使用 submitter 線程池異步提交剩余任務,不阻塞調用方 |
| 立即返回 | submitAll() 立即返回,上層可用 Futures.allAsList() 組合 |
核心知識點
1. ExecutorCompletionService
JDK提供的任務完成服務,核心能力是按完成順序獲取結果。
ExecutorService pool = Executors.newFixedThreadPool(10);
ExecutorCompletionService<Train> cs = new ExecutorCompletionService<>(pool);
// 提交任務
cs.submit(() -> queryTrain("G1"));
cs.submit(() -> queryTrain("G2"));
// 阻塞獲取最先完成的結果(而不是提交順序)
Future<Train> first = cs.take(); // 誰先完成就拿到誰
Future<Train> second = cs.take();原理:
┌─────────────────────────────────────────────┐
│ ExecutorCompletionService │
├─────────────────────────────────────────────┤
│ - executor: ExecutorService │
│ - completionQueue: BlockingQueue<Future> │
├─────────────────────────────────────────────┤
│ + submit(Callable): Future │
│ + take(): Future // 阻塞獲取已完成的 │
│ + poll(): Future // 非阻塞獲取 │
└─────────────────────────────────────────────┘
│
│ 任務完成時自動放入隊列
▼
┌─────────────────────────────────────────────┐
│ BlockingQueue<Future<V>> │
│ ┌──────┐ ┌──────┐ ┌──────┐ │
│ │Future│ │Future│ │Future│ ... │
│ └──────┘ └──────┘ └──────┘ │
│ 已完成 已完成 已完成 │
└─────────────────────────────────────────────┘關鍵點:
- 提交任務時包裝Future,添加完成回調
- 任務完成后自動將Future放入
completionQueue take()從隊列獲取,天然按完成順序
2. Guava ListenableFuture
標準Future的增強版,支持回調監(jiān)聽。
// 標準Future:只能輪詢或阻塞
Future<Train> future = executor.submit(() -> queryTrain("G1"));
Train result = future.get(); // 阻塞
// ListenableFuture:注冊回調
ListeningExecutorService executor = MoreExecutors.listeningDecorator(pool);
ListenableFuture<Train> future = executor.submit(() -> queryTrain("G1"));
// 異步回調(不阻塞)
future.addListener(() -> {
System.out.println("G1查詢完成!");
}, directExecutor());
// 或使用Futures工具類
Futures.addCallback(future, new FutureCallback<Train>() {
public void onSuccess(Train result) { ... }
public void onFailure(Throwable t) { ... }
}, executor);核心價值:
- 非阻塞:不需要輪詢
isDone()或阻塞get() - 組合能力:
Futures.allAsList(),transform(),catching()等 - 異常傳播:失敗會自動傳播到下游Future
3. SettableFuture
可手動設置結果的Future,類似CompletableFuture。
SettableFuture<Train> future = SettableFuture.create(); // 在其他線程設置結果 future.set(train); // 正常完成 future.setException(new TimeoutException()); // 異常完成 future.setFuture(otherFuture); // 鏈接另一個Future // 調用方正常使用 Train result = future.get();
4. 并發(fā)度控制核心思路
滑動窗口策略:
假設 parallelism=3, tasks=10(查詢10個車次)
初始狀態(tài):提交前3個車次到線程池
┌───┬───┬───┐ ┌───┬───┬───┬───┬───┬───┬───┐
│ G1│ G2│ G3│ │ G4│ G5│ G6│ G7│ G8│ G9│G10│
└───┴───┴───┘ └───┴───┴───┴───┴───┴───┴───┘
執(zhí)行中 等待中
G1 完成 → 立即提交 G4
┌───┬───┬───┐ ┌───┬───┬───┬───┬───┬───┐
│ │ G2│ G3│ │ G4│ G5│ G6│ G7│ G8│ G9│G10│
└───┴───┴───┘ └───┴───┴───┴───┴───┴───┴───┘
執(zhí)行中 等待中
G2 完成 → 立即提交 G5
...
維持窗口大小 = parallelism,任務完成后立即補充新任務實現(xiàn)要點:
- 使用
ExecutorCompletionService.take()等待任意任務完成 - 每完成一個,立即從待執(zhí)行隊列提交下一個
- 保證同時執(zhí)行的任務數 ≤ parallelism
使用場景
| 場景 | 說明 |
|---|---|
| 批量查詢接口 | 火車票、機票、酒店查詢,下游限流保護 |
| 批量緩存查詢 | Redis 連接池有限,控制并發(fā)數 |
| 批量 IO 操作 | 文件讀寫、數據庫查詢 |
| 爬蟲系統(tǒng) | 限制爬取速率,避免被封 |
總結
- ExecutorCompletionService -
take()按完成順序獲取,支持動態(tài)補充任務 - ListenableFuture + SettableFuture - 通過
linkFuture將內部Future鏈接到返回占位符 - 雙線程池設計 - pool執(zhí)行任務,submitter異步補充,不阻塞調用方
- 滑動窗口策略 - 初始提交N個,每完成一個立即補充下一個
- 索引嚴格對應 - result數量等于tasks數量,result[i]對應tasks[i]
核心優(yōu)勢:
- 精確控制并發(fā)度,避免下游限流
- 立即返回Future,支持流式處理
- 資源利用率高,短任務不等待長任務
- 代碼簡潔,無需手動管理信號量
- 代碼沒有實現(xiàn)中斷處理與取消傳播,感興趣的讀者可以自行實現(xiàn)(思路:linkFuture支持雙向取消,提交線程響應中斷+處理隊列毒丸,考慮以下幾種模式 fail-fast、fail-last、cancelN)
到此這篇關于Java滑動窗口實現(xiàn)線程池并發(fā)度控制的文章就介紹到這了,更多相關Java線程池并發(fā)度控制內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
SpringBoot自定義實現(xiàn)內容協(xié)商的三種策略
內容協(xié)商是HTTP協(xié)議中的一個重要概念,允許同一資源URL根據客戶端的偏好提供不同格式的表示,這篇文章主要介紹了SpringBoot自定義實現(xiàn)內容協(xié)商的三種策略,希望對大家有一定的幫助2025-04-04

