Java線程池高效并發(fā)編程實戰(zhàn)技巧
概要
因為面試中暴露出來的不足,所以寫一寫線程池,也算是復(fù)習(xí)一下。
什么是線程池
線程池是一種線程管理機制,它預(yù)先創(chuàng)建一定數(shù)量的線程并放入池中,當(dāng)需要執(zhí)行任務(wù)時,從池中獲取空閑線程來執(zhí)行任務(wù),任務(wù)完成后線程不銷毀而是返回池中等待下一次任務(wù)。
主要作用
1、降低資源消耗
避免頻繁創(chuàng)建和銷毀線程的開銷,重復(fù)利用已經(jīng)創(chuàng)建好了的線程。
2、提高響應(yīng)速度
任務(wù)到達時,無需額外創(chuàng)建線程即可運行
3、提高線程可管理性
統(tǒng)一管理線程資源,避免無限制創(chuàng)建線程導(dǎo)致系統(tǒng)崩潰,可以控制并發(fā)線程數(shù)量,避免過度競爭。
4、提供更強大的功能
- 定時執(zhí)行,周期執(zhí)行
- 任務(wù)隊列管理
- 拒絕策略
基礎(chǔ)線程池使用實例
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秒,任務(wù)隊列容量100
ThreadPoolExecutor executor = new ThreadPoolExecutor(
5, // corePoolSize
10, // maximumPoolSize
60, // keepAliveTime
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100), // 任務(wù)隊列
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy() // 拒絕策略
);
// 2. 提交任務(wù)
for (int i = 1; i <= 20; i++) {
int taskId = i;
executor.execute(() -> {
System.out.println("處理任務(wù)" + taskId +
", 線程: " + Thread.currentThread().getName());
try {
Thread.sleep(1000); // 模擬業(yè)務(wù)處理
} 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,任務(wù)隊列容量為100
2、提交任務(wù)
- 循環(huán)提交:代碼通過
for循環(huán)向線程池提交了 20 個任務(wù)。 - 變量捕獲:這里定義
int taskId = i;是因為在 Lambda 表達式內(nèi)部引用的外部變量必須是 final 或 effectively final(即不再改變)。直接用i會報錯,因為i在循環(huán)中一直在變。 - 非阻塞:
executor.execute()是異步的。這意味著主線程會瞬間跑完這個循環(huán),把 20 個任務(wù)丟進線程池的任務(wù)隊列,而不會等待任務(wù)執(zhí)行完。
3、關(guān)閉
線程池的工作流程
當(dāng)你調(diào)用 executor.execute() 時,內(nèi)部會發(fā)生以下邏輯:
- 核心線程(Core Threads):如果當(dāng)前運行的線程少于核心線程數(shù),直接創(chuàng)建新線程執(zhí)行。
- 任務(wù)隊列(Work Queue):如果核心線程滿了,任務(wù)會進入隊列排隊。
- 最大線程(Max Threads):如果隊列也滿了,且線程數(shù)少于最大線程數(shù),則創(chuàng)建非核心線程。
- 拒絕策略:如果全都滿了,就會觸發(fā)拒絕策略(Reject Policy)。
業(yè)務(wù)場景
訂單異步處理
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);
}
// 等待所有任務(wù)完成
return CompletableFuture.allOf(
futures.toArray(new CompletableFuture[0])
);
}
// 業(yè)務(wù)方法(模擬實現(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è)務(wù)邏輯,啥都不用管。但是缺點就是,你什么都不知道,無法等待任務(wù)完成,而且不知道任務(wù)會不會被線程池拒絕。
2、這么寫有什么好處呢?
答:
雖然你的邏輯內(nèi)部不產(chǎn)生結(jié)果(即 runAsync 的特性),但返回 CompletableFuture 有以下三個核心好處:
- 鏈?zhǔn)秸{(diào)用: 調(diào)用者可以寫
processOrderAsync(order).thenRun(() -> System.out.println("全部搞定"))。 - 異常處理: 調(diào)用者可以使用
.exceptionally()統(tǒng)一處理異步鏈路中的崩潰。 - 等待結(jié)束: 在單元測試或系統(tǒng)關(guān)閉前,可以調(diào)用
.join()確保任務(wù)執(zhí)行完了。
(關(guān)于鏈?zhǔ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ù)導(dǎo)出
import java.util.concurrent.*;
import java.util.List;
import java.io.File;
public class DataExportService {
// 專門用于導(dǎo)出任務(wù)的線程池
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() // 拒絕策略:丟棄最老任務(wù)
);
/**
* 異步導(dǎo)出Excel
*/
public CompletableFuture<File> exportExcelAsync(String exportId,
List<?> dataList) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("開始導(dǎo)出數(shù)據(jù),任務(wù)ID: " + exportId);
try {
// 模擬大數(shù)據(jù)量處理
File excelFile = generateExcelFile(dataList);
// 模擬上傳到云存儲
String url = uploadToCloudStorage(excelFile);
// 記錄導(dǎo)出日志
saveExportLog(exportId, url, "SUCCESS");
return excelFile;
} catch (Exception e) {
saveExportLog(exportId, null, "FAILED");
throw new RuntimeException("導(dǎo)出失敗", e);
}
}, exportExecutor);
}
/**
* 帶進度的數(shù)據(jù)導(dǎo)出
*/
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è)務(wù)方法(模擬實現(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("任務(wù)隊列大小: " + exportExecutor.getQueue().size());
System.out.println("已完成任務(wù)數(shù): " + exportExecutor.getCompletedTaskCount());
}
}說明
這里有點看不懂,先說一下吧,exportExcelAsync是最標(biāo)準(zhǔn)的異步流導(dǎo)出,直接調(diào)用線程池進行業(yè)務(wù)邏輯的調(diào)用并且返回結(jié)果,流程如下:
這是最基礎(chǔ)的異步流,采用了 CompletableFuture.supplyAsync。
執(zhí)行步驟:
- 提交任務(wù):將任務(wù)交給
exportExecutor處理。 - 生成文件:調(diào)用
generateExcelFile(模擬耗時操作)。 - 上傳云端:將生成的 File 上傳到 OSS 等存儲服務(wù)。
- 保存日志:無論成功還是失敗,都會記錄
saveExportLog。 - 返回結(jié)果:返回一個
File對象,調(diào)用者可以通過.get()或.thenAccept()獲取。
而exportWithProgress則是這樣子的
這是這段代碼的高級之處。它解決了大數(shù)據(jù)量導(dǎo)出時“用戶不知道還要等多久”的問題。
- 分批處理 (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塊中捕獲異常,并在失敗時記錄錯誤日志,確保即便導(dǎo)出崩了,系統(tǒng)也知道原因。 - 狀態(tài)監(jiān)控 (
printThreadPoolStatus):提供了一個監(jiān)控入口。在實際生產(chǎn)中,我們可以通過這個方法觀察隊列是否積壓,從而判斷是否需要增加核心線程數(shù)。
可以優(yōu)化的點:
- 拒絕策略的風(fēng)險:
DiscardOldestPolicy會讓某些用戶永遠等不到他們的文件(任務(wù)被悄悄丟棄了)。在金融或嚴(yán)肅業(yè)務(wù)中,通常改用CallerRunsPolicy(讓調(diào)用者自己執(zhí)行)或者自定義異常拋出。 - 內(nèi)存占用:
List<?> dataList如果非常大(比如百萬級),直接傳入方法可能會導(dǎo)致 OOM (內(nèi)存溢出)。通常建議傳入查詢條件,在異步線程里分頁從數(shù)據(jù)庫讀取。 - 線程中斷:
exportWithProgress里的Thread.sleep捕獲了中斷信號并重置了狀態(tài),這是非常專業(yè)的寫法,值得點贊。
也就是這里
// 模擬處理時間
try {
Thread.sleep(100);
} catch (InterruptedException e) {
// 就是這一句!重新設(shè)置中斷狀態(tài)
Thread.currentThread().interrupt();
}為什么說這行代碼“很專業(yè)”?
在 Java 并發(fā)編程中,這是一個非常容易被新手忽略的最佳實踐。
1. 中斷標(biāo)志位被“擦除”了
當(dāng)一個線程正在 sleep 時,如果外部調(diào)用了 thread.interrupt(),sleep 方法會立刻拋出 InterruptedException。 重點來了: 一旦拋出這個異常,JVM 會自動把該線程的“中斷標(biāo)志位”清除(改為 false)。
2. 如果不加這一句會發(fā)生什么?
如果你只是打印了日志,或者干脆 catch 塊里什么都不寫:
- 線程的中斷狀態(tài)丟失了。
- 上層代碼(或者線程池的后續(xù)邏輯)無法知道這個線程曾經(jīng)被要求停止。
- 這就像是有人按了“緊急停止”按鈕,結(jié)果系統(tǒng)捕捉到了信號但轉(zhuǎn)頭就給忘了,導(dǎo)致程序繼續(xù)盲目運行。
3. Thread.currentThread().interrupt() 的作用
這一行的意思是:“既然異常把中斷標(biāo)志位擦除了,那我就手動把它再設(shè)回 true。”
這樣做有幾個好處:
- 傳遞信號: 如果這個任務(wù)后續(xù)還有其他的檢查點(比如
Thread.currentThread().isInterrupted()),它能感知到中斷。 - 尊重規(guī)范: 讓線程池(
exportExecutor)或更高層的調(diào)用者能看到線程的中斷狀態(tài),從而決定是否回收線程或停止后續(xù)任務(wù)。
定時任務(wù)線程池
import java.util.concurrent.*;
import java.time.LocalDateTime;
public class ScheduledTaskService {
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(3);
/**
* 初始化定時任務(wù)
*/
public void initScheduledTasks() {
// 1. 每天凌晨執(zhí)行數(shù)據(jù)清理
scheduleDailyCleanup();
// 2. 每5分鐘執(zhí)行一次數(shù)據(jù)同步
schedulePeriodicSync();
// 3. 延遲執(zhí)行一次性任務(wù)
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í)行一次性任務(wù)
*/
private void scheduleOneTimeTask() {
scheduler.schedule(() -> {
System.out.println("執(zhí)行一次性任務(wù): " + LocalDateTime.now());
}, 10, TimeUnit.SECONDS);
}
/**
* 提交可取消的定時任務(wù)
*/
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è)務(wù)方法
private void cleanUpOldData() {
// 清理過期數(shù)據(jù)
}
private void syncDataWithExternalSystem() {
// 同步數(shù)據(jù)
}
public void shutdown() {
scheduler.shutdown();
}
}說明
這個之前做過,這里總結(jié)一下真正業(yè)務(wù)中會怎么做
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è)務(wù)邏輯
} catch (Exception e) {
log.error("清理失敗", e);
}
}
}2、分布式鎖
代碼在單機運行沒問題,但現(xiàn)代業(yè)務(wù)通常是 多實例部署。
- 痛點:如果部署了 3 個節(jié)點,凌晨 2 點時,3 個節(jié)點會同時跑清理任務(wù),可能導(dǎo)致數(shù)據(jù)庫死鎖或重復(fù)處理。
- 方案:使用 ShedLock 或 Redis 鎖,確保同一時間只有一個實例執(zhí)行。
(當(dāng)時的統(tǒng)計數(shù)據(jù)業(yè)務(wù)就是這么處理的)
@Scheduled(cron = "0 0 2 * * ?")
@SchedulerLock(name = "dataCleanupTask", lockAtMostFor = "10m", lockAtLeastFor = "1m")
public void scheduledTask() {
// 只有搶到鎖的機器才會執(zhí)行
}3、分布式任務(wù)調(diào)度平臺(XXL-JOB / Quartz)
在大型互聯(lián)網(wǎng)公司,定時任務(wù)通常是獨立于業(yè)務(wù)代碼進行管理的。最常用的方案是 XXL-JOB(國內(nèi)主流)或 Elastic-Job。
為什么業(yè)務(wù)開發(fā)喜歡用平臺?
- 可視化管理:不需要改代碼,在網(wǎng)頁上就能開關(guān)任務(wù)、修改執(zhí)行時間。
- 彈性調(diào)度:如果一臺服務(wù)器掛了,平臺會自動把任務(wù)調(diào)度到另一臺健康的服務(wù)器。
- 失敗告警:任務(wù)失敗了會自動發(fā)郵件/釘釘通知,還有重試機制。
- 執(zhí)行日志:平臺記錄了每次執(zhí)行的耗時、結(jié)果,方便排查。
總結(jié)
| 場景 | 推薦方案 |
| 本地小工具/單機腳本 | 維持你現(xiàn)在的 ScheduledExecutorService (最輕量) |
| 普通 Spring Boot 業(yè)務(wù) | @Scheduled + 配置文件 |
| 多臺服務(wù)器集群部署 | @Scheduled + ShedLock (最簡單有效) |
| 中大型分布式系統(tǒng) | XXL-JOB 或 Cloud Native CronJob (最專業(yè)) |
小結(jié)
對于線程池的用法做了一點小小的總結(jié),這是個開始。
到此這篇關(guān)于Java線程池高效并發(fā)編程實戰(zhàn)技巧的文章就介紹到這了,更多相關(guān)Java線程池并發(fā)編程內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
maven中自定義MavenArchetype的實現(xiàn)
Maven自身提供了許多Archetype來方便用戶創(chuàng)建Project,為了避免在創(chuàng)建project時重復(fù)的拷貝和修改,我們通過自定義Archetype來規(guī)范顯得還蠻有必要,下面就來介紹一下,感興趣的可以了解一下2025-01-01
Mybatis中自定義TypeHandler處理枚舉的示例代碼
typeHandler,是 MyBatis 中的一個接口,用于處理數(shù)據(jù)庫中的特定數(shù)據(jù)類型,下面簡單介紹創(chuàng)建自定義 typeHandler 來處理枚舉類型的示例,感興趣的朋友跟隨小編一起看看吧2024-01-01
java控制臺實現(xiàn)學(xué)生信息管理系統(tǒng)
這篇文章主要為大家詳細介紹了java控制臺實現(xiàn)學(xué)生信息管理系統(tǒng),文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下2022-02-02
Java swing實現(xiàn)支持錄音等功能的鋼琴程序
這篇文章主要為大家詳細介紹了Java swing實現(xiàn)鋼琴程序,支持錄音等功能的Java鋼琴源碼,具有一定的參考價值,感興趣的小伙伴們可以參考一下2017-06-06
idea快捷鍵生成getter和setter,有構(gòu)造參數(shù),無構(gòu)造參數(shù),重寫toString方式
這篇文章主要介紹了java之idea快捷鍵生成getter和setter,有構(gòu)造參數(shù),無構(gòu)造參數(shù),重寫toString方式,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教2023-11-11
Spring Cloud Config實現(xiàn)分布式配置中心
這篇文章主要介紹了Spring Cloud Config實現(xiàn)分布式配置中心,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2018-04-04

