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

Java線程池高效并發(fā)編程實戰(zhàn)技巧

 更新時間:2026年02月04日 09:36:41   作者:舊日之血_Hayter  
線程池是一種線程管理機制,它預先創(chuàng)建一定數(shù)量的線程并放入池中,當需要執(zhí)行任務時,從池中獲取空閑線程來執(zhí)行任務,任務完成后線程不銷毀而是返回池中等待下一次任務,這篇文章給大家介紹Java線程池高效并發(fā)編程實戰(zhàn)技巧,感興趣的朋友跟隨小編一起看看吧

概要

因為面試中暴露出來的不足,所以寫一寫線程池,也算是復習一下。

什么是線程池

線程池是一種線程管理機制,它預先創(chuàng)建一定數(shù)量的線程并放入池中,當需要執(zhí)行任務時,從池中獲取空閑線程來執(zhí)行任務,任務完成后線程不銷毀而是返回池中等待下一次任務。

主要作用

1、降低資源消耗

避免頻繁創(chuàng)建和銷毀線程的開銷,重復利用已經(jīng)創(chuàng)建好了的線程。

2、提高響應速度

任務到達時,無需額外創(chuàng)建線程即可運行

3、提高線程可管理性

統(tǒng)一管理線程資源,避免無限制創(chuàng)建線程導致系統(tǒng)崩潰,可以控制并發(fā)線程數(shù)量,避免過度競爭。

4、提供更強大的功能

  • 定時執(zhí)行,周期執(zhí)行
  • 任務隊列管理
  • 拒絕策略

基礎線程池使用實例

import java.util.concurrent.*;
import java.util.Random;
public class ThreadPoolDemo {
    public static void main(String[] args) {
        // 1. 創(chuàng)建線程池
        // 核心參數(shù):核心線程數(shù)5,最大線程數(shù)10,空閑時間60秒,任務隊列容量100
        ThreadPoolExecutor executor = new ThreadPoolExecutor(
            5,  // corePoolSize
            10, // maximumPoolSize
            60, // keepAliveTime
            TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(100), // 任務隊列
            Executors.defaultThreadFactory(),
            new ThreadPoolExecutor.AbortPolicy() // 拒絕策略
        );
        // 2. 提交任務
        for (int i = 1; i <= 20; i++) {
            int taskId = i;
            executor.execute(() -> {
                System.out.println("處理任務" + taskId + 
                    ", 線程: " + Thread.currentThread().getName());
                try {
                    Thread.sleep(1000); // 模擬業(yè)務處理
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            });
        }
        // 3. 優(yōu)雅關(guān)閉
        executor.shutdown();
        try {
            if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
                executor.shutdownNow();
            }
        } catch (InterruptedException e) {
            executor.shutdownNow();
        }
    }
}

如上所示,我們經(jīng)歷了:

1、創(chuàng)建線程池

其核心線程數(shù)為5,最大線程數(shù)為10,空閑時間60s,任務隊列容量為100

2、提交任務

  • 循環(huán)提交:代碼通過 for 循環(huán)向線程池提交了 20 個任務。
  • 變量捕獲:這里定義 int taskId = i; 是因為在 Lambda 表達式內(nèi)部引用的外部變量必須是 finaleffectively final(即不再改變)。直接用 i 會報錯,因為 i 在循環(huán)中一直在變。
  • 非阻塞executor.execute()異步的。這意味著主線程會瞬間跑完這個循環(huán),把 20 個任務丟進線程池的任務隊列,而不會等待任務執(zhí)行完。

3、關(guān)閉

線程池的工作流程

當你調(diào)用 executor.execute() 時,內(nèi)部會發(fā)生以下邏輯:

  • 核心線程(Core Threads):如果當前運行的線程少于核心線程數(shù),直接創(chuàng)建新線程執(zhí)行。
  • 任務隊列(Work Queue):如果核心線程滿了,任務會進入隊列排隊。
  • 最大線程(Max Threads):如果隊列也滿了,且線程數(shù)少于最大線程數(shù),則創(chuàng)建非核心線程。
  • 拒絕策略:如果全都滿了,就會觸發(fā)拒絕策略(Reject Policy)。

業(yè)務場景

訂單異步處理

import java.util.concurrent.*;
import java.util.List;
import java.util.ArrayList;
public class OrderProcessor {
    // 使用單例模式創(chuàng)建線程池
    private static final ThreadPoolExecutor orderExecutor = new ThreadPoolExecutor(
        3, 8, 30, TimeUnit.SECONDS,
        new ArrayBlockingQueue<>(1000),
        new ThreadFactory() {
            private int count = 0;
            @Override
            public Thread newThread(Runnable r) {
                Thread thread = new Thread(r, "order-process-" + (++count));
                thread.setDaemon(false);
                return thread;
            }
        },
        new ThreadPoolExecutor.CallerRunsPolicy() // 拒絕策略:調(diào)用者線程執(zhí)行
    );
    /**
     * 異步處理訂單
     */
    public CompletableFuture<Void> processOrderAsync(Order order) {
        return CompletableFuture.runAsync(() -> {
            try {
                // 1. 驗證訂單
                validateOrder(order);
                // 2. 扣減庫存
                reduceInventory(order);
                // 3. 生成發(fā)貨單
                generateShipping(order);
                // 4. 發(fā)送通知
                sendNotification(order);
                System.out.println("訂單處理完成: " + order.getId());
            } catch (Exception e) {
                // 記錄異常,進行補償
                handleOrderException(order, e);
            }
        }, orderExecutor);
    }
    /**
     * 批量處理訂單
     */
    public CompletableFuture<Void> batchProcessOrders(List<Order> orders) {
        List<CompletableFuture<Void>> futures = new ArrayList<>();
        for (Order order : orders) {
            CompletableFuture<Void> future = processOrderAsync(order);
            futures.add(future);
        }
        // 等待所有任務完成
        return CompletableFuture.allOf(
            futures.toArray(new CompletableFuture[0])
        );
    }
    // 業(yè)務方法(模擬實現(xiàn))
    private void validateOrder(Order order) {
        // 驗證邏輯
    }
    private void reduceInventory(Order order) {
        // 扣減庫存邏輯
    }
    private void generateShipping(Order order) {
        // 生成發(fā)貨單邏輯
    }
    private void sendNotification(Order order) {
        // 發(fā)送通知
    }
    private void handleOrderException(Order order, Exception e) {
        // 異常處理
    }
    // 優(yōu)雅關(guān)閉
    public void shutdown() {
        orderExecutor.shutdown();
    }
    // 訂單類
    static class Order {
        private String id;
        // 其他字段
        public String getId() { return id; }
    }
}

說明

1、為什么返回類型是CompletableFuture<Void>?

答:

其實你要是不想返回的話,直接void就行。線程自己處理業(yè)務邏輯,啥都不用管。但是缺點就是,你什么都不知道,無法等待任務完成,而且不知道任務會不會被線程池拒絕。

2、這么寫有什么好處呢?

答:

雖然你的邏輯內(nèi)部不產(chǎn)生結(jié)果(即 runAsync 的特性),但返回 CompletableFuture 有以下三個核心好處:

  • 鏈式調(diào)用: 調(diào)用者可以寫 processOrderAsync(order).thenRun(() -> System.out.println("全部搞定"))。
  • 異常處理: 調(diào)用者可以使用 .exceptionally() 統(tǒng)一處理異步鏈路中的崩潰。
  • 等待結(jié)束: 在單元測試或系統(tǒng)關(guān)閉前,可以調(diào)用 .join() 確保任務執(zhí)行完了。

(關(guān)于鏈式調(diào)用的問題,后面會新開一遍文章說一下,愛你。)

3、還有什么常見的返回類型嗎?

返回類型場景建議
CompletableFuture<Void>推薦。 異步執(zhí)行,不返回數(shù)據(jù),但允許調(diào)用者監(jiān)聽狀態(tài)。
void極致的“甩手掌柜”,調(diào)用方完全不關(guān)心后續(xù),代碼最簡。
CompletableFuture<T>異步執(zhí)行,且需要把處理后的結(jié)果傳回給調(diào)用方。

異步數(shù)據(jù)導出

import java.util.concurrent.*;
import java.util.List;
import java.io.File;
public class DataExportService {
    // 專門用于導出任務的線程池
    private static final ThreadPoolExecutor exportExecutor = new ThreadPoolExecutor(
        2, 4, 5, TimeUnit.MINUTES,
        new LinkedBlockingQueue<>(50),
        new ThreadFactory() {
            private int count = 0;
            @Override
            public Thread newThread(Runnable r) {
                Thread thread = new Thread(r, "export-thread-" + (++count));
                thread.setPriority(Thread.NORM_PRIORITY);
                return thread;
            }
        },
        new ThreadPoolExecutor.DiscardOldestPolicy() // 拒絕策略:丟棄最老任務
    );
    /**
     * 異步導出Excel
     */
    public CompletableFuture<File> exportExcelAsync(String exportId, 
                                                     List<?> dataList) {
        return CompletableFuture.supplyAsync(() -> {
            System.out.println("開始導出數(shù)據(jù),任務ID: " + exportId);
            try {
                // 模擬大數(shù)據(jù)量處理
                File excelFile = generateExcelFile(dataList);
                // 模擬上傳到云存儲
                String url = uploadToCloudStorage(excelFile);
                // 記錄導出日志
                saveExportLog(exportId, url, "SUCCESS");
                return excelFile;
            } catch (Exception e) {
                saveExportLog(exportId, null, "FAILED");
                throw new RuntimeException("導出失敗", e);
            }
        }, exportExecutor);
    }
    /**
     * 帶進度的數(shù)據(jù)導出
     */
    public CompletableFuture<File> exportWithProgress(String exportId, 
                                                      List<?> dataList,
                                                      ProgressCallback callback) {
        return CompletableFuture.supplyAsync(() -> {
            int total = dataList.size();
            int batchSize = 1000;
            int processed = 0;
            for (int i = 0; i < total; i += batchSize) {
                int end = Math.min(i + batchSize, total);
                List<?> batchData = dataList.subList(i, end);
                // 處理批次數(shù)據(jù)
                processBatchData(batchData);
                processed = end;
                float progress = (float) processed / total;
                // 回調(diào)更新進度
                if (callback != null) {
                    callback.onProgress(progress);
                }
                // 模擬處理時間
                try {
                    Thread.sleep(100);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            }
            return generateExcelFile(dataList);
        }, exportExecutor);
    }
    // 業(yè)務方法(模擬實現(xiàn))
    private File generateExcelFile(List<?> dataList) {
        // 生成Excel文件
        return new File("export.xlsx");
    }
    private String uploadToCloudStorage(File file) {
        // 上傳到云存儲
        return "https://oss.example.com/" + file.getName();
    }
    private void saveExportLog(String exportId, String url, String status) {
        // 保存日志
    }
    private void processBatchData(List<?> batchData) {
        // 處理批次數(shù)據(jù)
    }
    // 進度回調(diào)接口
    public interface ProgressCallback {
        void onProgress(float progress);
    }
    // 獲取線程池狀態(tài)
    public void printThreadPoolStatus() {
        System.out.println("核心線程數(shù): " + exportExecutor.getCorePoolSize());
        System.out.println("活動線程數(shù): " + exportExecutor.getActiveCount());
        System.out.println("任務隊列大小: " + exportExecutor.getQueue().size());
        System.out.println("已完成任務數(shù): " + exportExecutor.getCompletedTaskCount());
    }
}

說明

這里有點看不懂,先說一下吧,exportExcelAsync是最標準的異步流導出,直接調(diào)用線程池進行業(yè)務邏輯的調(diào)用并且返回結(jié)果,流程如下:

這是最基礎的異步流,采用了 CompletableFuture.supplyAsync。

執(zhí)行步驟:

  • 提交任務:將任務交給 exportExecutor 處理。
  • 生成文件:調(diào)用 generateExcelFile(模擬耗時操作)。
  • 上傳云端:將生成的 File 上傳到 OSS 等存儲服務。
  • 保存日志:無論成功還是失敗,都會記錄 saveExportLog。
  • 返回結(jié)果:返回一個 File 對象,調(diào)用者可以通過 .get().thenAccept() 獲取。

而exportWithProgress則是這樣子的

這是這段代碼的高級之處。它解決了大數(shù)據(jù)量導出時“用戶不知道還要等多久”的問題。

  • 分批處理 (Batching):它通過 for 循環(huán)和 subList 將原始數(shù)據(jù)切分成每 1000 條一組。
  • 進度計算:每次處理完一批,計算 processed / total 的百分比。
  • 回調(diào)機制 (ProgressCallback):每完成一個批次,就調(diào)用一次 callback.onProgress(progress)。
    • 注意: 這個回調(diào)通常會連接到 WebSocket 或 Redis,從而讓前端頁面能實時顯示進度條。
  • 模擬延遲Thread.sleep(100) 是為了模擬真實處理數(shù)據(jù)的耗時,防止瞬時完成看不出進度效果。

兩者都用了try catch來保證健壯性

  • 異常處理:在 try-catch 塊中捕獲異常,并在失敗時記錄錯誤日志,確保即便導出崩了,系統(tǒng)也知道原因。
  • 狀態(tài)監(jiān)控 (printThreadPoolStatus):提供了一個監(jiān)控入口。在實際生產(chǎn)中,我們可以通過這個方法觀察隊列是否積壓,從而判斷是否需要增加核心線程數(shù)。

可以優(yōu)化的點:

  • 拒絕策略的風險DiscardOldestPolicy 會讓某些用戶永遠等不到他們的文件(任務被悄悄丟棄了)。在金融或嚴肅業(yè)務中,通常改用 CallerRunsPolicy(讓調(diào)用者自己執(zhí)行)或者自定義異常拋出。
  • 內(nèi)存占用List<?> dataList 如果非常大(比如百萬級),直接傳入方法可能會導致 OOM (內(nèi)存溢出)。通常建議傳入查詢條件,在異步線程里分頁從數(shù)據(jù)庫讀取。
  • 線程中斷exportWithProgress 里的 Thread.sleep 捕獲了中斷信號并重置了狀態(tài),這是非常專業(yè)的寫法,值得點贊。

也就是這里

// 模擬處理時間
try {
    Thread.sleep(100);
} catch (InterruptedException e) {
    // 就是這一句!重新設置中斷狀態(tài)
    Thread.currentThread().interrupt();
}

為什么說這行代碼“很專業(yè)”?

在 Java 并發(fā)編程中,這是一個非常容易被新手忽略的最佳實踐。

1. 中斷標志位被“擦除”了

當一個線程正在 sleep 時,如果外部調(diào)用了 thread.interrupt(),sleep 方法會立刻拋出 InterruptedException重點來了: 一旦拋出這個異常,JVM 會自動把該線程的“中斷標志位”清除(改為 false)。

2. 如果不加這一句會發(fā)生什么?

如果你只是打印了日志,或者干脆 catch 塊里什么都不寫:

  • 線程的中斷狀態(tài)丟失了。
  • 上層代碼(或者線程池的后續(xù)邏輯)無法知道這個線程曾經(jīng)被要求停止。
  • 這就像是有人按了“緊急停止”按鈕,結(jié)果系統(tǒng)捕捉到了信號但轉(zhuǎn)頭就給忘了,導致程序繼續(xù)盲目運行。

3. Thread.currentThread().interrupt() 的作用

這一行的意思是:“既然異常把中斷標志位擦除了,那我就手動把它再設回 true。”

這樣做有幾個好處:

  • 傳遞信號: 如果這個任務后續(xù)還有其他的檢查點(比如 Thread.currentThread().isInterrupted()),它能感知到中斷。
  • 尊重規(guī)范: 讓線程池(exportExecutor)或更高層的調(diào)用者能看到線程的中斷狀態(tài),從而決定是否回收線程或停止后續(xù)任務。

定時任務線程池

import java.util.concurrent.*;
import java.time.LocalDateTime;
public class ScheduledTaskService {
    private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(3);
    /**
     * 初始化定時任務
     */
    public void initScheduledTasks() {
        // 1. 每天凌晨執(zhí)行數(shù)據(jù)清理
        scheduleDailyCleanup();
        // 2. 每5分鐘執(zhí)行一次數(shù)據(jù)同步
        schedulePeriodicSync();
        // 3. 延遲執(zhí)行一次性任務
        scheduleOneTimeTask();
    }
    /**
     * 每天凌晨2點執(zhí)行數(shù)據(jù)清理
     */
    private void scheduleDailyCleanup() {
        long initialDelay = calculateInitialDelay(2, 0); // 凌晨2點
        long period = 24 * 60 * 60; // 24小時
        scheduler.scheduleAtFixedRate(() -> {
            try {
                System.out.println("開始數(shù)據(jù)清理: " + LocalDateTime.now());
                cleanUpOldData();
                System.out.println("數(shù)據(jù)清理完成: " + LocalDateTime.now());
            } catch (Exception e) {
                System.err.println("數(shù)據(jù)清理失敗: " + e.getMessage());
            }
        }, initialDelay, period, TimeUnit.SECONDS);
    }
    /**
     * 每5分鐘執(zhí)行數(shù)據(jù)同步
     */
    private void schedulePeriodicSync() {
        scheduler.scheduleWithFixedDelay(() -> {
            try {
                syncDataWithExternalSystem();
            } catch (Exception e) {
                // 記錄異常,下次繼續(xù)執(zhí)行
                System.err.println("數(shù)據(jù)同步失敗: " + e.getMessage());
            }
        }, 0, 5, TimeUnit.MINUTES);
    }
    /**
     * 延遲10秒執(zhí)行一次性任務
     */
    private void scheduleOneTimeTask() {
        scheduler.schedule(() -> {
            System.out.println("執(zhí)行一次性任務: " + LocalDateTime.now());
        }, 10, TimeUnit.SECONDS);
    }
    /**
     * 提交可取消的定時任務
     */
    public ScheduledFuture<?> submitCancellableTask(Runnable task, 
                                                    long initialDelay, 
                                                    long period, 
                                                    TimeUnit unit) {
        return scheduler.scheduleAtFixedRate(task, initialDelay, period, unit);
    }
    // 工具方法:計算到指定時間的延遲
    private long calculateInitialDelay(int targetHour, int targetMinute) {
        LocalDateTime now = LocalDateTime.now();
        LocalDateTime targetTime = now.withHour(targetHour)
                                     .withMinute(targetMinute)
                                     .withSecond(0);
        if (now.isAfter(targetTime)) {
            targetTime = targetTime.plusDays(1);
        }
        return java.time.Duration.between(now, targetTime).getSeconds();
    }
    // 業(yè)務方法
    private void cleanUpOldData() {
        // 清理過期數(shù)據(jù)
    }
    private void syncDataWithExternalSystem() {
        // 同步數(shù)據(jù)
    }
    public void shutdown() {
        scheduler.shutdown();
    }
}

說明

這個之前做過,這里總結(jié)一下真正業(yè)務中會怎么做

1、Spring的用法

如果項目是 Spring Boot,通常不會手動去 new ScheduledExecutorService。我們會利用 Spring 封裝好的注解,配合配置文件。

  • 優(yōu)點:代碼極其簡潔,支持 Cron 表達式。
  • 企業(yè)級改法:將時間配置寫在 application.yml 或配置中心(Apollo/Nacos)。
@Component
@Slf4j
public class DataCleanupTask {
    // 從配置文件讀取 Cron 表達式,例如:0 0 2 * * ? (每天凌晨2點)
    @Scheduled(cron = "${task.cleanup.cron}")
    public void dailyCleanup() {
        log.info("開始數(shù)據(jù)清理...");
        try {
            // 業(yè)務邏輯
        } catch (Exception e) {
            log.error("清理失敗", e);
        }
    }
}

2、分布式鎖

代碼在單機運行沒問題,但現(xiàn)代業(yè)務通常是 多實例部署。

  • 痛點:如果部署了 3 個節(jié)點,凌晨 2 點時,3 個節(jié)點會同時跑清理任務,可能導致數(shù)據(jù)庫死鎖或重復處理。
  • 方案:使用 ShedLock 或 Redis 鎖,確保同一時間只有一個實例執(zhí)行。

(當時的統(tǒng)計數(shù)據(jù)業(yè)務就是這么處理的)

@Scheduled(cron = "0 0 2 * * ?")
@SchedulerLock(name = "dataCleanupTask", lockAtMostFor = "10m", lockAtLeastFor = "1m")
public void scheduledTask() {
    // 只有搶到鎖的機器才會執(zhí)行
}

3、分布式任務調(diào)度平臺(XXL-JOB / Quartz)

在大型互聯(lián)網(wǎng)公司,定時任務通常是獨立于業(yè)務代碼進行管理的。最常用的方案是 XXL-JOB(國內(nèi)主流)或 Elastic-Job。

為什么業(yè)務開發(fā)喜歡用平臺?

  • 可視化管理:不需要改代碼,在網(wǎng)頁上就能開關(guān)任務、修改執(zhí)行時間。
  • 彈性調(diào)度:如果一臺服務器掛了,平臺會自動把任務調(diào)度到另一臺健康的服務器。
  • 失敗告警:任務失敗了會自動發(fā)郵件/釘釘通知,還有重試機制。
  • 執(zhí)行日志:平臺記錄了每次執(zhí)行的耗時、結(jié)果,方便排查。

總結(jié)

場景推薦方案
本地小工具/單機腳本維持你現(xiàn)在的 ScheduledExecutorService (最輕量)
普通 Spring Boot 業(yè)務@Scheduled + 配置文件
多臺服務器集群部署@Scheduled + ShedLock (最簡單有效)
中大型分布式系統(tǒng)XXL-JOBCloud Native CronJob (最專業(yè))

小結(jié)

對于線程池的用法做了一點小小的總結(jié),這是個開始。

到此這篇關(guān)于Java線程池高效并發(fā)編程實戰(zhàn)技巧的文章就介紹到這了,更多相關(guān)Java線程池并發(fā)編程內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評論

杨浦区| 江川县| 伊春市| 称多县| 闵行区| 阳东县| 方城县| 淳化县| 寿阳县| 财经| 广宁县| 石城县| 灵武市| 永泰县| 菏泽市| 永城市| 思南县| 交城县| 西和县| 永康市| 上栗县| 武隆县| 东台市| 灯塔市| 桃园县| 瓮安县| 林周县| 儋州市| 景洪市| 灵寿县| 天镇县| 安图县| 米易县| 元阳县| 双城市| 商南县| 凤山市| 南平市| 沭阳县| 衢州市| 龙江县|