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

Java滑動窗口實現(xiàn)線程池并發(fā)度控制的方法詳解

 更新時間:2026年04月13日 10:57:11   作者:樺說編程  
Java并發(fā)線程池是一種用于管理和復用線程的機制,它簡化了線程的生命周期管理,提高了資源利用率和程序響應速度,這篇文章主要介紹了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ù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • Java并發(fā)編程變量可見性避免指令重排使用詳解

    Java并發(fā)編程變量可見性避免指令重排使用詳解

    這篇文章主要為大家介紹了Java并發(fā)編程變量可見性避免指令重排使用詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2022-11-11
  • 使用jasypt在springboot中加密敏感信息

    使用jasypt在springboot中加密敏感信息

    Jasypt是一個簡化Java加密操作的輕量級庫,支持與Spring Boot深度集成,通過ENC(密文)格式實現(xiàn)配置文件的自動加解密,本文就來介紹一下springboot jasypt 加密敏感,感興趣的可以了解一下
    2026-02-02
  • 詳解java裝飾模式(Decorator Pattern)

    詳解java裝飾模式(Decorator Pattern)

    這篇文章主要為大家詳細介紹了java裝飾模式Decorator Pattern,這種類型的設計模式屬于結構型模式,它是作為現(xiàn)有的類的一個包裝,對裝飾器模式感興趣的小伙伴們可以參考一下
    2016-04-04
  • 解決spring data redis的那些坑

    解決spring data redis的那些坑

    這篇文章主要介紹了spring data redis的那些坑,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Java中Iterator迭代器的使用詳解

    Java中Iterator迭代器的使用詳解

    在程序開發(fā)中,經常需要遍歷集合中的所有元素。針對這種需求,JDK專門提供了一個接口java.util.Iterator。本文就來詳細說說Iterator迭代器的使用,感興趣的可以了解一下
    2022-10-10
  • Java模擬棧和隊列數據結構的基本示例講解

    Java模擬棧和隊列數據結構的基本示例講解

    這篇文章主要介紹了Java模擬棧和隊列數據結構的基本示例,棧的后進先出和隊列的先進先出是數據結構中最基礎的知識,本文則又對Java實現(xiàn)棧和隊列結構的方法進行了細分,需要的朋友可以參考下
    2016-04-04
  • Java并發(fā)編程之Fork/Join框架詳解

    Java并發(fā)編程之Fork/Join框架詳解

    這篇文章主要介紹了Java并發(fā)編程之Fork/Join框架詳解,Fork/Join框架是Java7提供的一個用于并行執(zhí)行任務的框架,是一個把大任務分割成若干個小任務,最終匯總每個小任務結果后得到大任務結果的框架,需要的朋友可以參考下
    2023-12-12
  • SpringBoot自定義實現(xiàn)內容協(xié)商的三種策略

    SpringBoot自定義實現(xiàn)內容協(xié)商的三種策略

    內容協(xié)商是HTTP協(xié)議中的一個重要概念,允許同一資源URL根據客戶端的偏好提供不同格式的表示,這篇文章主要介紹了SpringBoot自定義實現(xiàn)內容協(xié)商的三種策略,希望對大家有一定的幫助
    2025-04-04
  • Java中InputSteam怎么轉String

    Java中InputSteam怎么轉String

    面了一位實習生,叫他給我說一下怎么把InputStream轉換為String,這種常規(guī)的操作,他竟然都沒有用過我準備結合工作經驗,整理匯集出了InputStream 到String 轉換的十八般武藝,助大家闖蕩 Java 江湖一臂之力,需要的朋友可以參考下
    2021-06-06
  • Java如何替換jar中的class文件

    Java如何替換jar中的class文件

    在調整java代碼過程中會遇到需要改jar包中的class文件的情況,改了如何替換呢?下面小編給大家分享java替換jar中的class文件的操作方法,感興趣的朋友跟隨小編一起看看吧
    2024-02-02

最新評論

兰溪市| 阆中市| 葫芦岛市| 唐海县| 隆林| 长丰县| 林甸县| 雷山县| 临夏县| 从化市| 五家渠市| 兴文县| 浪卡子县| 赤水市| 贞丰县| 抚宁县| 新巴尔虎右旗| 洮南市| 水城县| 大竹县| 江北区| 阿克| 瑞昌市| 新泰市| 仁怀市| 通山县| 侯马市| 英超| 宣城市| 融水| 琼结县| 罗甸县| 连平县| 五寨县| 柳林县| 桃园县| 资兴市| 万荣县| 荥经县| 渝北区| 抚宁县|