Java大數(shù)據(jù)量異步處理方案
一、為什么需要異步
當(dāng)一次操作需要處理大量數(shù)據(jù)(如插入12萬(wàn)條記錄到數(shù)據(jù)庫(kù)),如果同步執(zhí)行:
- 用戶等待時(shí)間過(guò)長(zhǎng)(可能幾十秒到幾分鐘)
- HTTP 連接可能超時(shí)
- 服務(wù)器線程被長(zhǎng)時(shí)間占用,影響其他請(qǐng)求
解決思路:先快速響應(yīng)用戶"任務(wù)已提交",再在后臺(tái)異步完成耗時(shí)操作。
二、方案對(duì)比
2.1 線程池(ThreadPoolExecutor)
原理:在 JVM 內(nèi)部維護(hù)一組工作線程,將任務(wù)提交到隊(duì)列中由這些線程異步執(zhí)行。
優(yōu)點(diǎn):
- 零依賴:不需要額外中間件
- 低延遲:任務(wù)提交后立即被線程拾取執(zhí)行
- 簡(jiǎn)單直接:代碼量少,調(diào)試方便
- 適合單服務(wù)內(nèi)部的異步任務(wù)
缺點(diǎn):
- 不可靠:JVM 重啟或崩潰時(shí),隊(duì)列中未執(zhí)行的任務(wù)丟失
- 不可分布式:只能在本機(jī)執(zhí)行,無(wú)法分發(fā)到其他節(jié)點(diǎn)
- 隊(duì)列有限:隊(duì)列滿了要么阻塞、要么拒絕
- 不可觀測(cè):沒(méi)有天然的任務(wù)狀態(tài)追蹤、重試機(jī)制
適用場(chǎng)景:
- 數(shù)據(jù)丟失可接受(重新導(dǎo)入即可)
- 單實(shí)例部署或任務(wù)不需要跨實(shí)例分發(fā)
- 對(duì)實(shí)時(shí)性要求高(毫秒級(jí)開(kāi)始執(zhí)行)
2.2 消息隊(duì)列(RabbitMQ / RocketMQ / Kafka)
原理:將任務(wù)以消息形式發(fā)送到 Broker,消費(fèi)者從 Broker 拉取消息執(zhí)行。
優(yōu)點(diǎn):
- 高可靠:消息持久化,Broker 宕機(jī)恢復(fù)后消息不丟
- 可分布式:多個(gè)消費(fèi)者實(shí)例分擔(dān)負(fù)載
- 削峰填谷:突發(fā)流量堆積在隊(duì)列中,消費(fèi)者按自身速度處理
- 天然可重試:消費(fèi)失敗可重新入隊(duì)
- 可觀測(cè):有管理控制臺(tái)查看隊(duì)列積壓、消費(fèi)進(jìn)度
缺點(diǎn):
- 引入外部依賴(Broker 部署、運(yùn)維、配置)
- 增加系統(tǒng)復(fù)雜度(消息序列化、冪等性、順序性)
- 延遲略高(網(wǎng)絡(luò)往返 + Broker 中轉(zhuǎn))
- 調(diào)試?yán)щy(異步鏈路追蹤)
適用場(chǎng)景:
- 任務(wù)不能丟失,必須保證執(zhí)行
- 多實(shí)例部署需要負(fù)載分發(fā)
- 需要削峰(如秒殺、批量任務(wù)集中提交)
- 需要跨服務(wù)通信
2.3 Spring @Async
原理:通過(guò)注解標(biāo)記方法為異步,Spring 使用內(nèi)部線程池執(zhí)行。
優(yōu)點(diǎn):
- 極簡(jiǎn):加個(gè)注解就行
- 聲明式:不需要手動(dòng)管理線程池
缺點(diǎn):
- 底層還是線程池,有線程池的所有缺點(diǎn)
- 默認(rèn)線程池配置不合理(SimpleAsyncTaskExecutor 每次創(chuàng)建新線程)
- 事務(wù)傳播復(fù)雜:異步方法中的事務(wù)與調(diào)用方獨(dú)立
- 自調(diào)用失效:同一個(gè)類內(nèi)部調(diào)用 @Async 方法不會(huì)異步(代理問(wèn)題)
適用場(chǎng)景:
- 簡(jiǎn)單異步任務(wù)
- 對(duì)線程池參數(shù)不需要精細(xì)控制
2.4 對(duì)比表
| 維度 | 線程池 | 消息隊(duì)列 | @Async |
|---|---|---|---|
| 可靠性 | 低(JVM 重啟丟失) | 高(消息持久化) | 低 |
| 分布式 | ? | ? | ? |
| 外部依賴 | 無(wú) | 需要 Broker | 無(wú) |
| 延遲 | 極低(微秒級(jí)) | 低(毫秒級(jí)) | 極低 |
| 削峰能力 | 有限(隊(duì)列大小) | 強(qiáng)(Broker 容量) | 有限 |
| 代碼復(fù)雜度 | 中 | 高 | 低 |
| 可觀測(cè)性 | 弱 | 強(qiáng) | 弱 |
| 重試機(jī)制 | 需自行實(shí)現(xiàn) | 內(nèi)置 | 需自行實(shí)現(xiàn) |
三、線程池核心知識(shí)
3.1 ThreadPoolExecutor 七大參數(shù)
new ThreadPoolExecutor(
corePoolSize, // 核心線程數(shù):始終存活的線程
maximumPoolSize, // 最大線程數(shù):隊(duì)列滿了之后擴(kuò)展到的上限
keepAliveTime, // 空閑線程存活時(shí)間
timeUnit, // 時(shí)間單位
workQueue, // 任務(wù)隊(duì)列
threadFactory, // 線程工廠(自定義線程名稱)
rejectedHandler // 拒絕策略
);
3.2 任務(wù)提交執(zhí)行流程
提交任務(wù) ├── 當(dāng)前線程數(shù) < corePoolSize → 創(chuàng)建新核心線程執(zhí)行 ├── 當(dāng)前線程數(shù) >= corePoolSize → 放入 workQueue ├── workQueue 已滿 且 當(dāng)前線程數(shù) < maximumPoolSize → 創(chuàng)建非核心線程執(zhí)行 └── workQueue 已滿 且 當(dāng)前線程數(shù) >= maximumPoolSize → 執(zhí)行拒絕策略
3.3 四種拒絕策略
| 策略 | 行為 | 適用場(chǎng)景 |
|---|---|---|
| AbortPolicy | 拋出 RejectedExecutionException | 不允許丟任務(wù),調(diào)用方需感知 |
| CallerRunsPolicy | 由提交任務(wù)的線程自己執(zhí)行 | 不丟任務(wù),自動(dòng)降級(jí)為同步 |
| DiscardPolicy | 靜默丟棄 | 允許丟失 |
| DiscardOldestPolicy | 丟棄隊(duì)列中最老的任務(wù) | 只關(guān)心最新任務(wù) |
3.4 常見(jiàn)隊(duì)列選擇
| 隊(duì)列類型 | 特點(diǎn) |
|---|---|
| ArrayBlockingQueue | 有界,背壓明確 |
| LinkedBlockingQueue | 可有界可無(wú)界,無(wú)界時(shí)可能 OOM |
| SynchronousQueue | 零容量,直接交接(用于 CachedThreadPool) |
3.5 參數(shù)設(shè)計(jì)經(jīng)驗(yàn)
CPU 密集型任務(wù)(計(jì)算、排序):
- corePoolSize = CPU 核心數(shù) + 1
- 隊(duì)列可以短一些
IO 密集型任務(wù)(數(shù)據(jù)庫(kù)寫(xiě)入、網(wǎng)絡(luò)調(diào)用):
- corePoolSize = CPU 核心數(shù) × 2 或更高
- 線程大部分時(shí)間在等待 IO,可以多一些
批量導(dǎo)入場(chǎng)景(大量數(shù)據(jù)庫(kù)寫(xiě)入):
- corePoolSize 不需要太大(4~8),避免數(shù)據(jù)庫(kù)連接池被打滿
- 隊(duì)列適當(dāng)大(16~32),允許少量堆積
- 拒絕策略用 CallerRunsPolicy,保證不丟任務(wù)
四、完整示例:基于線程池的異步批量數(shù)據(jù)導(dǎo)入
以下是一個(gè)示例,展示"同步校驗(yàn) + 異步批量插入"模式。
4.1 線程池配置
package com.example.config;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* 異步導(dǎo)入線程池配置.
*/
@Configuration
public class AsyncImportThreadPoolConfig {
@Bean(name = "importThreadPool", destroyMethod = "shutdown")
public ThreadPoolExecutor importThreadPool() {
return new ThreadPoolExecutor(
4, // 核心線程數(shù)
8, // 最大線程數(shù)
60, TimeUnit.SECONDS, // 空閑線程存活60秒
new ArrayBlockingQueue<>(16), // 有界隊(duì)列,最多堆積16個(gè)任務(wù)
new ImportThreadFactory(), // 自定義線程工廠
new ThreadPoolExecutor.CallerRunsPolicy() // 隊(duì)列滿時(shí)由調(diào)用線程執(zhí)行
);
}
static class ImportThreadFactory implements ThreadFactory {
private final AtomicInteger counter = new AtomicInteger(1);
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, "import-worker-" + counter.getAndIncrement());
t.setDaemon(false); // 非守護(hù)線程,確保任務(wù)執(zhí)行完
return t;
}
}
}
4.2 Service 接口
package com.example.service;
import java.util.List;
/**
* 批量導(dǎo)入服務(wù)接口.
*/
public interface BatchImportService {
/**
* 導(dǎo)入數(shù)據(jù):同步校驗(yàn) + 異步入庫(kù).
*
* @param rawDataList 原始數(shù)據(jù)列表(已從文件中解析出來(lái))
* @param operatorId 操作人ID
* @return 導(dǎo)入結(jié)果提示
*/
String importData(List<RawData> rawDataList, String operatorId);
}
4.3 Service 實(shí)現(xiàn)
package com.example.service.impl;
import com.example.entity.ImportRecord;
import com.example.mapper.ImportRecordMapper;
import com.example.service.BatchImportService;
import jakarta.annotation.Resource;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ThreadPoolExecutor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
@Slf4j
@Service
public class BatchImportServiceImpl implements BatchImportService {
private static final int BATCH_SIZE = 2000;
@Resource
@Qualifier("importThreadPool")
private ThreadPoolExecutor threadPool;
@Resource
private ImportRecordMapper importRecordMapper;
@Override
public String importData(List<RawData> rawDataList, String operatorId) {
// ========== 第一步:同步校驗(yàn)(在請(qǐng)求線程中執(zhí)行) ==========
for (int i = 0; i < rawDataList.size(); i++) {
RawData raw = rawDataList.get(i);
String error = validate(raw);
if (error != null) {
// 遇到第一條錯(cuò)誤立即中斷,同步返回給前端
throw new RuntimeException("第" + (i + 1) + "行:" + error);
}
}
// ========== 第二步:數(shù)據(jù)轉(zhuǎn)換 ==========
List<ImportRecord> recordList = new ArrayList<>(rawDataList.size());
for (RawData raw : rawDataList) {
ImportRecord record = convertToEntity(raw, operatorId);
recordList.add(record);
}
// ========== 第三步:異步批量插入(提交到線程池) ==========
// 注意:這里的 recordList 對(duì)象引用傳遞給了異步線程
// 確保主線程之后不再修改這個(gè)列表
threadPool.execute(() -> {
try {
long start = System.currentTimeMillis();
int total = recordList.size();
for (int i = 0; i < total; i += BATCH_SIZE) {
int end = Math.min(i + BATCH_SIZE, total);
List<ImportRecord> batch = recordList.subList(i, end);
importRecordMapper.batchInsert(batch);
}
long cost = System.currentTimeMillis() - start;
log.info("異步導(dǎo)入完成,共{}條,耗時(shí){}ms", total, cost);
} catch (Exception e) {
log.error("異步導(dǎo)入失敗", e);
// 可選:更新主表狀態(tài)為"導(dǎo)入失敗"
}
});
// ========== 第四步:同步返回成功提示 ==========
return "導(dǎo)入任務(wù)已提交,共" + recordList.size() + "條數(shù)據(jù)正在后臺(tái)處理";
}
/**
* 校驗(yàn)單條數(shù)據(jù).
* 返回 null 表示通過(guò),返回錯(cuò)誤信息表示失敗.
*/
private String validate(RawData raw) {
if (raw.getAmount() == null) {
return "數(shù)量不能為空";
}
if (raw.getAmount() <= 0 || raw.getAmount() > 999999) {
return "數(shù)量必須為大于0的正整數(shù),最多六位";
}
return null;
}
/**
* 原始數(shù)據(jù)轉(zhuǎn)換為實(shí)體.
*/
private ImportRecord convertToEntity(RawData raw, String operatorId) {
ImportRecord record = new ImportRecord();
record.setCode(raw.getCode());
record.setName(raw.getName());
record.setAmount(raw.getAmount());
record.setOperatorId(operatorId);
return record;
}
}
4.4 執(zhí)行時(shí)序
請(qǐng)求線程 線程池工作線程 │ │ │── 解析文件 ──→ │ │── 逐行校驗(yàn) ──→ │ │ (校驗(yàn)不過(guò)直接返回錯(cuò)誤) │ │── 轉(zhuǎn)換數(shù)據(jù) ──→ │ │── threadPool.execute(task) ──→ │ │ │── 批量INSERT第1批(2000條) │← 返回"導(dǎo)入任務(wù)已提交" ── │── 批量INSERT第2批(2000條) │ │── ... │ (HTTP響應(yīng)已返回給前端) │── 批量INSERT第N批 │ │── 記錄日志"導(dǎo)入完成"
五、線程池方案的注意事項(xiàng)
5.1 線程安全
提交給線程池的數(shù)據(jù)(如 recordList)在主線程返回后不能再修改。示例中使用的是 ArrayList,提交后主線程不再操作它,所以安全。如果有并發(fā)修改風(fēng)險(xiǎn),應(yīng)使用 Collections.unmodifiableList() 或復(fù)制一份。
5.2 事務(wù)邊界
異步線程中的數(shù)據(jù)庫(kù)操作有獨(dú)立的事務(wù)上下文。如果需要在主表保存后、從表插入中途失敗時(shí)回滾主表,需要額外的補(bǔ)償邏輯(如更新主表狀態(tài)為"導(dǎo)入失敗")。
5.3 優(yōu)雅停機(jī)
Spring Boot 配置 server.shutdown=graceful 后,停機(jī)時(shí)會(huì)等待請(qǐng)求處理完成。但線程池中的任務(wù)默認(rèn)不被等待。配置 destroyMethod = "shutdown" 可以讓 Spring 容器銷(xiāo)毀 Bean 時(shí)調(diào)用 shutdown(),等待正在執(zhí)行的任務(wù)完成(但隊(duì)列中等待的任務(wù)不會(huì)執(zhí)行)。
如果要確保隊(duì)列中的任務(wù)也執(zhí)行完:
@PreDestroy
public void destroy() {
threadPool.shutdown();
try {
if (!threadPool.awaitTermination(60, TimeUnit.SECONDS)) {
threadPool.shutdownNow();
}
} catch (InterruptedException e) {
threadPool.shutdownNow();
}
}
5.4 監(jiān)控
線程池沒(méi)有內(nèi)置的管理界面。建議通過(guò)定時(shí)任務(wù)或 Actuator 暴露:
log.info("線程池狀態(tài) - 活躍:{}, 隊(duì)列積壓:{}, 已完成:{}",
threadPool.getActiveCount(),
threadPool.getQueue().size(),
threadPool.getCompletedTaskCount());
六、什么時(shí)候該用消息隊(duì)列替代線程池
| 信號(hào) | 建議 |
|---|---|
| 服務(wù)多實(shí)例部署,需要負(fù)載均衡消費(fèi) | 用消息隊(duì)列 |
| 任務(wù)絕對(duì)不能丟失(如金融交易) | 用消息隊(duì)列 |
| 需要延時(shí)執(zhí)行或定時(shí)重試 | 用消息隊(duì)列 |
| 任務(wù)量突增需要削峰 | 用消息隊(duì)列 |
| 單實(shí)例、任務(wù)可重試(如重新導(dǎo)入) | 線程池足夠 |
| 對(duì)延遲敏感(需要立即開(kāi)始執(zhí)行) | 線程池更合適 |
| 不想引入外部依賴 | 線程池 |
以上就是Java大數(shù)據(jù)量異步處理方案的詳細(xì)內(nèi)容,更多關(guān)于Java大數(shù)據(jù)量異步處理的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
一文學(xué)會(huì)處理SpringBoot統(tǒng)一返回格式
這篇文章主要介紹了一文學(xué)會(huì)處理SpringBoot統(tǒng)一返回格式,文章圍繞主題展開(kāi)詳細(xì)的內(nèi)容介紹,具有一定的參考價(jià)值,需要的小伙伴可以參考一下2022-08-08
全鏈路監(jiān)控平臺(tái)Pinpoint?SkyWalking?Zipkin選型對(duì)比
這篇文章主要為大家介紹了全鏈路監(jiān)控平臺(tái)Pinpoint?SkyWalking?Zipkin實(shí)現(xiàn)的選型對(duì)比,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步2022-03-03
JavaWeb實(shí)現(xiàn)文件上傳與下載實(shí)例詳解
在Web應(yīng)用程序開(kāi)發(fā)中,文件上傳與下載功能是非常常用的功能,下面通過(guò)本文給大家介紹JavaWeb實(shí)現(xiàn)文件上傳與下載實(shí)例詳解,對(duì)javaweb文件上傳下載相關(guān)知識(shí)感興趣的朋友一起學(xué)習(xí)吧2016-02-02
通過(guò)openOffice將office文件轉(zhuǎn)成pdf
這篇文章主要介紹了通過(guò)openOffice將office文件轉(zhuǎn)成pdf,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-11-11

