Excel 大数据量异步导入导出:Handler 注册机制 + 流式处理完整架构
Excel大数据量异步导入导出架构设计
引言
在企业级应用中,Excel导入导出是高频需求。当数据量达到数万甚至数十万行时,传统的同步处理方式会导致请求超时、内存溢出、用户体验差等问题。异步导入导出架构通过任务队列、流式处理、Handler注册机制等设计,实现了大数据量Excel的稳定高效处理。
本文将深入分析大数据量Excel处理的架构方案,探讨异步导出任务队列、导入校验框架、Handler注册机制解耦等核心设计。
一、传统方案的痛点
1.1 同步导出的问题
flowchart TD
subgraph 同步导出
A1[用户点击导出] --> A2[查询全部数据]
A2 --> A3[全量加载到内存]
A3 --> A4[生成Excel文件]
A4 --> A5[返回文件流]
A5 --> A6[浏览器下载]
A3 -.->|10万行数据| A7[内存占用: 2GB+]
A4 -.->|生成耗时| A8[响应时间: 30s+]
A5 -.->|超时风险| A9[网关504超时]
end
style A7 fill:#f44336,color:#fff
style A8 fill:#f44336,color:#fff
style A9 fill:#f44336,color:#fff
1.2 同步导入的问题
flowchart TD
subgraph 同步导入
B1[用户上传文件] --> B2[全量解析Excel]
B2 --> B3[逐行校验]
B3 --> B4[批量写入数据库]
B4 --> B5[返回导入结果]
B2 -.->|10万行文件| B6[解析耗时: 20s+]
B3 -.->|校验失败| B7[全部回滚/部分成功?]
B4 -.->|大批量INSERT| B8[数据库锁表]
end
style B6 fill:#f44336,color:#fff
style B7 fill:#f44336,color:#fff
style B8 fill:#f44336,color:#fff
二、异步导入导出整体架构
2.1 架构全景图
flowchart TD
subgraph 前端交互层
A1[导出按钮] --> A2[创建异步任务]
A3[导入上传] --> A4[创建异步任务]
A5[任务中心] --> A6[轮询任务状态]
A6 --> A7[下载文件/查看结果]
end
subgraph 任务管理层
A2 --> B1[TaskManager]
A4 --> B1
B1 --> B2[任务记录存DB]
B2 --> B3[提交到线程池]
end
subgraph Handler处理层
B3 --> C1{Handler路由}
C1 -->|user-export| C2[UserExportHandler]
C3 -->|user-import| C3[UserImportHandler]
C1 -->|order-export| C4[OrderExportHandler]
end
subgraph 数据处理层
C2 --> D1[分页查询+流式写入]
C3 --> D2[流式读取+批量校验]
D1 --> D3[生成Excel文件]
D2 --> D4[批量写入数据库]
end
subgraph 通知层
D3 --> E1[更新任务状态为完成]
D4 --> E1
E1 --> E2[WebSocket通知前端]
end
style B1 fill:#4CAF50,color:#fff
style C1 fill:#2196F3,color:#fff
style E2 fill:#FF9800,color:#fff
2.2 核心数据模型
/**
* 异步任务记录
*/
@Data
@TableName("sys_async_task")
public class AsyncTask {
@TableId(type = IdType.ASSIGN_ID)
private Long id;
/** 任务类型: EXPORT / IMPORT */
private String taskType;
/** 业务类型: user / order / product */
private String businessType;
/** 任务状态: PENDING / RUNNING / SUCCESS / FAIL */
private String status;
/** 任务参数(JSON格式) */
private String params;
/** 结果文件路径 */
private String filePath;
/** 总记录数 */
private Integer totalCount;
/** 成功记录数 */
private Integer successCount;
/** 失败记录数 */
private Integer failCount;
/** 错误信息 */
private String errorMsg;
/** 创建人 */
private Long createBy;
/** 创建时间 */
private LocalDateTime createTime;
/** 完成时间 */
private LocalDateTime finishTime;
}
三、Handler注册机制
3.1 Handler接口设计
/**
* 导出Handler接口
* 每种业务类型的导出逻辑实现此接口
*/
public interface ExportHandler<T> {
/**
* 获取Handler标识(对应businessType)
*/
String getBusinessType();
/**
* 分页查询数据
* @param params 查询参数
* @param pageNum 页码
* @param pageSize 每页大小
* @return 数据列表
*/
List<T> queryData(Map<String, Object> params, int pageNum, int pageSize);
/**
* 获取表头定义
*/
List<ExcelColumn> getColumns();
/**
* 获取总记录数
*/
long getTotalCount(Map<String, Object> params);
}
/**
* 导入Handler接口
* 每种业务类型的导入逻辑实现此接口
*/
public interface ImportHandler<T> {
/**
* 获取Handler标识
*/
String getBusinessType();
/**
* 将Excel行数据转换为实体
*/
T convertToEntity(Map<String, String> rowData);
/**
* 校验单条数据
* @return 校验错误信息,null表示通过
*/
String validate(T entity);
/**
* 批量保存数据
*/
void batchSave(List<T> entities);
}
3.2 Handler注册中心
/**
* Handler注册中心
* 通过Spring自动注册机制,将所有Handler实例注册到Map中
* 实现业务逻辑与调度框架的解耦
*/
@Component
@Slf4j
public class HandlerRegistry {
/** 导出Handler注册表 */
private final Map<String, ExportHandler<?>> exportHandlers = new ConcurrentHashMap<>();
/** 导入Handler注册表 */
private final Map<String, ImportHandler<?>> importHandlers = new ConcurrentHashMap<>();
/**
* 自动注册所有ExportHandler
*/
@Autowired
public void registerExportHandlers(List<ExportHandler<?>> handlers) {
handlers.forEach(handler -> {
exportHandlers.put(handler.getBusinessType(), handler);
log.info("注册导出Handler: businessType={}, handler={}",
handler.getBusinessType(), handler.getClass().getSimpleName());
});
}
/**
* 自动注册所有ImportHandler
*/
@Autowired
public void registerImportHandlers(List<ImportHandler<?>> handlers) {
handlers.forEach(handler -> {
importHandlers.put(handler.getBusinessType(), handler);
log.info("注册导入Handler: businessType={}, handler={}",
handler.getBusinessType(), handler.getClass().getSimpleName());
});
}
/**
* 获取导出Handler
*/
public ExportHandler<?> getExportHandler(String businessType) {
ExportHandler<?> handler = exportHandlers.get(businessType);
if (handler == null) {
throw new IllegalArgumentException("未注册的导出Handler: " + businessType);
}
return handler;
}
/**
* 获取导入Handler
*/
public ImportHandler<?> getImportHandler(String businessType) {
ImportHandler<?> handler = importHandlers.get(businessType);
if (handler == null) {
throw new IllegalArgumentException("未注册的导入Handler: " + businessType);
}
return handler;
}
}
3.3 Handler注册流程
flowchart TD
A[Spring容器启动] --> B[扫描所有ExportHandler实现]
B --> C[扫描所有ImportHandler实现]
C --> D[HandlerRegistry.registerExportHandlers]
D --> E[HandlerRegistry.registerImportHandlers]
E --> F[构建Handler映射表]
subgraph 映射表
F --> G["exportHandlers: {user → UserExportHandler, order → OrderExportHandler}"]
F --> H["importHandlers: {user → UserImportHandler, order → OrderImportHandler}"]
end
I[新增业务模块] --> J[实现ExportHandler/ImportHandler]
J --> K[添加@Component注解]
K --> L[自动注册到Registry]
L --> M[无需修改框架代码]
style F fill:#4CAF50,color:#fff
style M fill:#2196F3,color:#fff
四、异步导出实现
4.1 导出任务管理器
/**
* 异步导出任务管理器
*/
@Service
@Slf4j
@RequiredArgsConstructor
public class ExportTaskManager {
private final HandlerRegistry handlerRegistry;
private final AsyncTaskMapper taskMapper;
/** 导出专用线程池 */
private final ThreadPoolExecutor exportExecutor = new ThreadPoolExecutor(
2, 4, 60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100),
new ThreadFactoryBuilder().setNameFormat("export-pool-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
/**
* 创建并提交导出任务
*/
@Transactional
public Long submitExportTask(String businessType, Map<String, Object> params, Long userId) {
// 1. 创建任务记录
AsyncTask task = new AsyncTask();
task.setTaskType("EXPORT");
task.setBusinessType(businessType);
task.setParams(JSON.toJSONString(params));
task.setStatus("PENDING");
task.setCreateBy(userId);
task.setCreateTime(LocalDateTime.now());
taskMapper.insert(task);
// 2. 异步执行导出
exportExecutor.execute(() -> executeExport(task.getId(), businessType, params));
return task.getId();
}
/**
* 执行导出逻辑
*/
private void executeExport(Long taskId, String businessType, Map<String, Object> params) {
// 更新任务状态为运行中
updateTaskStatus(taskId, "RUNNING");
try {
ExportHandler<?> handler = handlerRegistry.getExportHandler(businessType);
long totalCount = handler.getTotalCount(params);
// 分页查询 + 流式写入Excel
String filePath = doExport(handler, params, totalCount);
// 更新任务状态为成功
AsyncTask update = new AsyncTask();
update.setId(taskId);
update.setStatus("SUCCESS");
update.setFilePath(filePath);
update.setTotalCount((int) totalCount);
update.setSuccessCount((int) totalCount);
update.setFinishTime(LocalDateTime.now());
taskMapper.updateById(update);
} catch (Exception e) {
log.error("导出任务执行失败: taskId={}", taskId, e);
updateTaskStatus(taskId, "FAIL", e.getMessage());
}
}
/**
* 分页查询 + 流式写入Excel
* 避免全量数据加载到内存
*/
private <T> String doExport(ExportHandler<T> handler, Map<String, Object> params, long total) {
String filePath = generateFilePath(handler.getBusinessType());
int pageSize = 5000; // 每页查询5000条
int totalPages = (int) Math.ceil((double) total / pageSize);
try (ExcelWriter writer = new ExcelWriter(filePath, handler.getColumns())) {
for (int page = 1; page <= totalPages; page++) {
List<T> data = handler.queryData(params, page, pageSize);
writer.writeRows(data);
log.debug("导出进度: {}/{}", page, totalPages);
}
}
return filePath;
}
}
4.2 流式Excel写入器
/**
* 流式Excel写入器
* 基于Apache POI的SXSSFWorkbook实现
* 每写入一定行数后自动刷新到磁盘,控制内存占用
*/
public class ExcelWriter implements AutoCloseable {
private final SXSSFWorkbook workbook;
private final Sheet sheet;
private final List<ExcelColumn> columns;
private int currentRow = 0;
/** 内存中保留的行数,超出后写入临时文件 */
private static final int ROW_ACCESS_WINDOW = 1000;
public ExcelWriter(String filePath, List<ExcelColumn> columns) throws IOException {
this.workbook = new SXSSFWorkbook(ROW_ACCESS_WINDOW);
this.sheet = workbook.createSheet();
this.columns = columns;
// 写入表头
writeHeader();
}
/**
* 写入表头行
*/
private void writeHeader() {
Row headerRow = sheet.createRow(currentRow++);
for (int i = 0; i < columns.size(); i++) {
Cell cell = headerRow.createCell(i);
cell.setCellValue(columns.get(i).getTitle());
// 表头样式
CellStyle style = workbook.createCellStyle();
Font font = workbook.createFont();
font.setBold(true);
style.setFont(font);
cell.setCellStyle(style);
}
}
/**
* 写入数据行
*/
public <T> void writeRows(List<T> dataList) {
for (T data : dataList) {
Row row = sheet.createRow(currentRow++);
Map<String, Object> fieldValues = toFieldMap(data);
for (int i = 0; i < columns.size(); i++) {
Cell cell = row.createCell(i);
Object value = fieldValues.get(columns.get(i).getField());
setCellValue(cell, value);
}
}
}
/**
* 将数据写入文件并关闭资源
*/
@Override
public void close() throws IOException {
try (FileOutputStream fos = new FileOutputStream(filePath)) {
workbook.write(fos);
}
// 清理临时文件
workbook.dispose();
}
}
五、异步导入实现
5.1 导入校验框架
flowchart TD
A[上传Excel文件] --> B[创建导入任务]
B --> C[流式读取Excel]
C --> D[逐行数据转换]
D --> E[逐行校验]
E --> F{校验结果}
F -->|通过| G[加入成功列表]
F -->|失败| H[记录错误信息]
G --> I{达到批次大小?}
I -->|是| J[批量写入数据库]
I -->|否| C
H --> C
C --> K{所有行处理完毕?}
K -->|否| C
K -->|是| L[写入剩余数据]
L --> M[生成导入结果报告]
M --> N[更新任务状态]
style E fill:#4CAF50,color:#fff
style J fill:#2196F3,color:#fff
style M fill:#FF9800,color:#fff
5.2 导入任务管理器
/**
* 异步导入任务管理器
*/
@Service
@Slf4j
@RequiredArgsConstructor
public class ImportTaskManager {
private final HandlerRegistry handlerRegistry;
private final AsyncTaskMapper taskMapper;
/** 导入专用线程池 */
private final ThreadPoolExecutor importExecutor = new ThreadPoolExecutor(
2, 4, 60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100),
new ThreadFactoryBuilder().setNameFormat("import-pool-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
/** 批量写入的批次大小 */
private static final int BATCH_SIZE = 500;
/**
* 创建并提交导入任务
*/
@Transactional
public Long submitImportTask(String businessType, String filePath, Long userId) {
AsyncTask task = new AsyncTask();
task.setTaskType("IMPORT");
task.setBusinessType(businessType);
task.setFilePath(filePath);
task.setStatus("PENDING");
task.setCreateBy(userId);
task.setCreateTime(LocalDateTime.now());
taskMapper.insert(task);
importExecutor.execute(() -> executeImport(task.getId(), businessType, filePath));
return task.getId();
}
/**
* 执行导入逻辑
*/
@SuppressWarnings("unchecked")
private void executeImport(Long taskId, String businessType, String filePath) {
updateTaskStatus(taskId, "RUNNING");
int totalCount = 0;
int successCount = 0;
int failCount = 0;
List<ImportError> errors = new ArrayList<>();
try {
ImportHandler<Object> handler =
(ImportHandler<Object>) handlerRegistry.getImportHandler(businessType);
List<Object> batchList = new ArrayList<>(BATCH_SIZE);
// 流式读取Excel
try (ExcelReader reader = new ExcelReader(filePath)) {
while (reader.hasNext()) {
Map<String, String> rowData = reader.nextRow();
totalCount++;
try {
// 数据转换
Object entity = handler.convertToEntity(rowData);
// 数据校验
String errorMsg = handler.validate(entity);
if (errorMsg != null) {
failCount++;
errors.add(new ImportError(totalCount, errorMsg, rowData.toString()));
continue;
}
// 加入批量列表
batchList.add(entity);
successCount++;
// 达到批次大小,批量写入
if (batchList.size() >= BATCH_SIZE) {
handler.batchSave(new ArrayList<>(batchList));
batchList.clear();
}
} catch (Exception e) {
failCount++;
errors.add(new ImportError(totalCount, e.getMessage(), rowData.toString()));
}
}
}
// 写入剩余数据
if (!batchList.isEmpty()) {
handler.batchSave(batchList);
}
// 更新任务结果
AsyncTask update = new AsyncTask();
update.setId(taskId);
update.setStatus("SUCCESS");
update.setTotalCount(totalCount);
update.setSuccessCount(successCount);
update.setFailCount(failCount);
update.setErrorMsg(errors.size() > 0 ? JSON.toJSONString(errors) : null);
update.setFinishTime(LocalDateTime.now());
taskMapper.updateById(update);
} catch (Exception e) {
log.error("导入任务执行失败: taskId={}", taskId, e);
updateTaskStatus(taskId, "FAIL", e.getMessage());
}
}
}
5.3 流式Excel读取器
/**
* 流式Excel读取器
* 基于Apache POI的SAX模式,逐行解析,内存占用恒定
*/
public class ExcelReader implements AutoCloseable, Iterable<Map<String, String>> {
private final OPCPackage pkg;
private final XSSFReader reader;
private final List<String> headers;
private Iterator<Map<String, String>> rowIterator;
public ExcelReader(String filePath) throws Exception {
this.pkg = OPCPackage.open(new File(filePath));
this.reader = new XSSFReader(pkg);
this.headers = parseHeaders();
}
/**
* 解析表头
*/
private List<String> parseHeaders() throws Exception {
// 读取第一行作为表头
List<String> result = new ArrayList<>();
// ... SAX解析逻辑
return result;
}
public boolean hasNext() {
return rowIterator != null && rowIterator.hasNext();
}
public Map<String, String> nextRow() {
return rowIterator.next();
}
@Override
public void close() throws IOException {
pkg.close();
}
}
六、前端交互设计
6.1 任务状态轮询
/**
* 异步任务轮询composable
*/
export function useAsyncTask() {
const taskStatus = ref<string>('')
const taskProgress = ref(0)
/**
* 提交导出任务并轮询状态
*/
const submitExport = async (businessType: string, params: Record<string, any>) => {
// 1. 创建任务
const { data: taskId } = await createExportTask(businessType, params)
taskStatus.value = 'PENDING'
// 2. 轮询任务状态
const pollTimer = setInterval(async () => {
const { data: task } = await getTaskStatus(taskId)
taskStatus.value = task.status
if (task.totalCount > 0) {
taskProgress.value = Math.round(
(task.successCount + task.failCount) / task.totalCount * 100
)
}
// 任务完成,停止轮询
if (['SUCCESS', 'FAIL'].includes(task.status)) {
clearInterval(pollTimer)
if (task.status === 'SUCCESS' && task.filePath) {
// 触发文件下载
downloadFile(task.filePath)
}
}
}, 2000) // 每2秒轮询一次
}
return { taskStatus, taskProgress, submitExport }
}
6.2 完整交互流程
sequenceDiagram
participant User as 用户
participant Frontend as 前端
participant API as 后端API
participant Pool as 线程池
participant Handler as ExportHandler
participant DB as 数据库
User->>Frontend: 点击导出
Frontend->>API: POST /api/export/create
API->>DB: 插入任务记录(PENDING)
API-->>Frontend: 返回taskId
API->>Pool: 提交异步任务
Frontend->>Frontend: 开始轮询(2s间隔)
Frontend->>API: GET /api/export/status/{taskId}
API-->>Frontend: {status: RUNNING, progress: 30%}
Pool->>Handler: 执行导出逻辑
Handler->>DB: 分页查询数据
Handler->>Handler: 流式写入Excel
Handler-->>Pool: 导出完成
Pool->>DB: 更新任务(SUCCESS)
Frontend->>API: GET /api/export/status/{taskId}
API-->>Frontend: {status: SUCCESS, filePath: "xxx"}
Frontend->>API: GET /api/export/download/{taskId}
API-->>User: 文件下载
七、性能优化策略
7.1 关键优化点
| 优化点 | 实现方式 | 效果 |
|---|---|---|
| 流式写入 | SXSSFWorkbook | 内存占用恒定,不受数据量影响 |
| 流式读取 | SAX模式解析 | 内存占用恒定,不受文件大小影响 |
| 分页查询 | 每次查询5000条 | 避免单次查询数据量过大 |
| 批量写入 | 每500条批量INSERT | 减少数据库交互次数 |
| 线程池隔离 | 导入/导出独立线程池 | 避免相互影响 |
| 任务超时 | 最大执行时间限制 | 防止任务无限运行 |
7.2 内存优化对比
flowchart LR
subgraph 传统方式
A1[10万行数据] --> A2[全量加载到内存]
A2 --> A3[内存峰值: 2GB+]
end
subgraph 流式处理
B1[10万行数据] --> B2[分页5000条加载]
B2 --> B3[内存峰值: ~200MB]
end
style A3 fill:#f44336,color:#fff
style B3 fill:#4CAF50,color:#fff
结论与建议
核心架构设计要点
- Handler注册机制:通过接口+自动注册实现业务逻辑与框架解耦,新增业务只需实现Handler接口
- 异步任务队列:线程池异步执行,前端轮询状态,避免请求超时
- 流式处理:导入使用SAX解析,导出使用SXSSFWorkbook,内存占用恒定
- 批量操作:导入批量写入数据库,导出分页查询数据,减少交互次数
- 校验框架:逐行校验+错误收集,支持部分成功导入
最佳实践建议
- 文件存储:导出文件存对象存储(MinIO/OSS),避免本地磁盘溢出
- 任务清理:定时清理过期的任务记录和文件(如7天后自动删除)
- 并发控制:限制同一用户同时执行的导入导出任务数
- 进度反馈:通过WebSocket实时推送任务进度,替代轮询方案
- 错误报告:导入失败时生成错误报告Excel,标注失败原因,方便用户修正