Java線程池拒絕策略原理及任務(wù)不丟失方案總結(jié)(最近實(shí)踐)
一、線程池拒絕策略的核心機(jī)制
Java線程池(ThreadPoolExecutor)的拒絕策略在以下條件下觸發(fā):
- 線程池已滿:
- 活躍線程數(shù) ≥
maximumPoolSize。 - 任務(wù)隊(duì)列(
workQueue)已滿(若為有界隊(duì)列)。
- 活躍線程數(shù) ≥
- 線程池關(guān)閉:
- 調(diào)用
shutdown()后提交新任務(wù)。
- 調(diào)用
觸發(fā)流程:
- 提交任務(wù)時(shí),線程池通過(guò)
execute()方法檢查狀態(tài)。 - 若線程池?zé)o法接受任務(wù)(如滿載或關(guān)閉),調(diào)用
RejectedExecutionHandler.rejectedExecution()。
二、四種內(nèi)置拒絕策略及適用場(chǎng)景
| 策略 | 行為 | 適用場(chǎng)景 | 風(fēng)險(xiǎn) |
|---|---|---|---|
| AbortPolicy | 拋出RejectedExecutionException | 關(guān)鍵任務(wù)(如支付) | 調(diào)用方需處理異常 |
| CallerRunsPolicy | 由提交任務(wù)的線程直接執(zhí)行任務(wù) | 非關(guān)鍵但需保證執(zhí)行(如日志上報(bào)) | 可能阻塞調(diào)用線程 |
| DiscardPolicy | 靜默丟棄任務(wù) | 可丟失任務(wù)(如監(jiān)控?cái)?shù)據(jù)) | 任務(wù)丟失風(fēng)險(xiǎn) |
| DiscardOldestPolicy | 丟棄隊(duì)列中最舊任務(wù),重試提交新任務(wù) | 實(shí)時(shí)性要求高(如股票行情) | 可能丟失重要任務(wù) |
三、確保任務(wù)不丟失的解決方案
1. 自定義拒絕策略 + 持久化存儲(chǔ)
核心思想:將拒絕的任務(wù)保存到外部存儲(chǔ)(數(shù)據(jù)庫(kù)/消息隊(duì)列/文件),后續(xù)通過(guò)重試機(jī)制恢復(fù)執(zhí)行。
實(shí)現(xiàn)步驟:
定義持久化任務(wù)實(shí)體:
@Data
public class PersistedTask {
private String id;
private String taskData; // 序列化后的任務(wù)
private int retryCount;
private LocalDateTime createTime;
}自定義拒絕策略:
public class PersistenceRejectPolicy implements RejectedExecutionHandler {
private final TaskRepository taskRepository;
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
try {
String taskData = serializeTask(r);
taskRepository.save(new PersistedTask(UUID.randomUUID().toString(), taskData, 0));
} catch (Exception e) {
log.error("持久化任務(wù)失敗", e);
}
}
private String serializeTask(Runnable r) {
return new Gson().toJson(r);
}
}定時(shí)任務(wù)重試:
@Scheduled(fixedRate = 5000)
public void retryRejectedTasks() {
List<PersistedTask> tasks = taskRepository.findPendingTasks();
for (PersistedTask task : tasks) {
try {
Runnable r = deserializeTask(task.getTaskData());
executor.execute(r);
taskRepository.markAsCompleted(task.getId());
} catch (Exception e) {
if (task.getRetryCount() >= 3) {
taskRepository.markAsFailed(task.getId());
} else {
taskRepository.incrementRetry(task.getId());
}
}
}
}2. 結(jié)合消息隊(duì)列的異步處理
優(yōu)勢(shì):解耦任務(wù)提交與執(zhí)行,利用消息隊(duì)列的持久化能力。
實(shí)現(xiàn):
public class MqRejectPolicy implements RejectedExecutionHandler {
private final KafkaTemplate<String, String> kafkaTemplate;
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
kafkaTemplate.send("rejected-tasks", serializeTask(r));
}
}3. 線程池參數(shù)優(yōu)化
- 隊(duì)列選擇:
- 有界隊(duì)列(如
ArrayBlockingQueue):需合理設(shè)置容量(例如queueCapacity = maxThreads * 2)。 - 優(yōu)先級(jí)隊(duì)列(
PriorityBlockingQueue):適合需要排序的任務(wù)。
- 有界隊(duì)列(如
- 線程數(shù)配置:
- CPU密集型:
corePoolSize = CPU核心數(shù)。 - IO密集型:
corePoolSize = 2 * CPU核心數(shù)。
- CPU密集型:
4. 優(yōu)雅關(guān)閉與任務(wù)完整性
public void shutdownGracefully(ThreadPoolExecutor executor) {
executor.shutdown(); // 拒絕新任務(wù)
try {
if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
List<Runnable> pendingTasks = executor.shutdownNow(); // 嘗試停止正在執(zhí)行的任務(wù)
savePendingTasks(pendingTasks); // 持久化未執(zhí)行任務(wù)
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}四、關(guān)鍵場(chǎng)景實(shí)踐
場(chǎng)景1:高并發(fā)訂單處理
- 策略:
AbortPolicy+ 數(shù)據(jù)庫(kù)持久化 + 熔斷機(jī)制。 - 配置:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
50, 200, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000),
new PersistenceRejectPolicy(taskRepository)
);
- 降級(jí):當(dāng)數(shù)據(jù)庫(kù)連接池滿時(shí),觸發(fā)熔斷直接丟棄非關(guān)鍵訂單。
場(chǎng)景2:實(shí)時(shí)數(shù)據(jù)分析
- 策略:
DiscardOldestPolicy+ Kafka緩沖。 - 配置:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
10, 50, 30, TimeUnit.SECONDS,
new PriorityBlockingQueue<>(100),
new MqRejectPolicy(kafkaTemplate)
);
五、風(fēng)險(xiǎn)與優(yōu)化建議
- 持久化性能瓶頸:
- 解決方案:批量插入數(shù)據(jù)庫(kù),或使用Redis List暫存任務(wù)。
- 任務(wù)重復(fù)執(zhí)行:
- 解決方案:為任務(wù)添加唯一ID,執(zhí)行前檢查是否已處理。
- 內(nèi)存泄漏:
- 解決方案:定期清理
EmergencyQueue中的積壓任務(wù)。
- 解決方案:定期清理
- 監(jiān)控缺失:
- 解決方案:通過(guò)Micrometer暴露以下指標(biāo):
threadpool.rejected.count:拒絕任務(wù)數(shù)。threadpool.queue.size:隊(duì)列堆積情況。
- 解決方案:通過(guò)Micrometer暴露以下指標(biāo):
總結(jié)
確保任務(wù)不丟失的核心在于:
- 拒絕策略選擇:根據(jù)業(yè)務(wù)容忍度選擇內(nèi)置策略或自定義持久化方案。
- 持久化設(shè)計(jì):結(jié)合數(shù)據(jù)庫(kù)/消息隊(duì)列存儲(chǔ)拒絕任務(wù)。
- 重試機(jī)制:通過(guò)定時(shí)任務(wù)或消費(fèi)者恢復(fù)任務(wù)。
- 系統(tǒng)優(yōu)化:合理配置線程池參數(shù),配合監(jiān)控與降級(jí)策略。
最佳實(shí)踐:在金融、電商等關(guān)鍵系統(tǒng)中,推薦自定義持久化策略 + 數(shù)據(jù)庫(kù)/Kafka + 重試機(jī)制;在日志上報(bào)等非關(guān)鍵場(chǎng)景,可使用CallerRunsPolicy + 本地緩存平衡性能與可靠性。
到此這篇關(guān)于Java線程池拒絕策略原理及任務(wù)不丟失方案總結(jié)(最近實(shí)踐)的文章就介紹到這了,更多相關(guān)java線程池拒絕策略內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Eureka源碼閱讀Client啟動(dòng)入口注冊(cè)續(xù)約及定時(shí)任務(wù)
這篇文章主要為大家介紹了Eureka源碼閱讀Client啟動(dòng)入口注冊(cè)續(xù)約及定時(shí)任務(wù)示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2022-10-10
基于Jackson實(shí)現(xiàn)API接口數(shù)據(jù)脫敏的示例詳解
用戶的一些敏感數(shù)據(jù),例如手機(jī)號(hào)、郵箱、身份證等信息,在數(shù)據(jù)庫(kù)以明文存儲(chǔ),但在接口返回?cái)?shù)據(jù)給瀏覽器(或三方客戶端)時(shí),希望對(duì)這些敏感數(shù)據(jù)進(jìn)行脫敏,所以本文就給大家介紹以惡如何利用Jackson實(shí)現(xiàn)API接口數(shù)據(jù)脫敏,需要的朋友可以參考下2023-08-08
IDEA中的Run/Debug Configurations各項(xiàng)解讀
這篇文章主要介紹了IDEA中的Run/Debug Configurations各項(xiàng)解讀,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-09-09
深入理解Java設(shè)計(jì)模式之訪問(wèn)者模式
這篇文章主要介紹了JAVA設(shè)計(jì)模式之訪問(wèn)者模式的的相關(guān)資料,文中示例代碼非常詳細(xì),供大家參考和學(xué)習(xí),感興趣的朋友可以了解2021-11-11
Java中單元測(cè)試框架JUnit知識(shí)點(diǎn)整理
在Java開(kāi)發(fā)中JUnit是最常用的單元測(cè)試框架之一,編寫(xiě)JUnit測(cè)試的目的是確保代碼的正確性、可維護(hù)性和可擴(kuò)展性,這篇文章主要介紹了Java中單元測(cè)試框架JUnit知識(shí)點(diǎn)整理的相關(guān)資料,需要的朋友可以參考下2025-08-08
Springboot項(xiàng)目長(zhǎng)時(shí)間不進(jìn)行接口操作,提示HikariPool-1警告的解決
這篇文章主要介紹了Springboot項(xiàng)目長(zhǎng)時(shí)間不進(jìn)行接口操作,提示HikariPool-1警告的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-12-12

