面对30万数据导出挑战,如何实现从系统崩溃到高效导出的技术演进

问题背景:一个看似简单的需求

最近接到一个业务需求:需要导出载体管理系统的30万条数据到Excel文件。这个看似简单的需求,却引发了系统性能危机:

  • 初始方案​:直接查询30万数据 → 内存溢出,系统崩溃
  • 第二次尝试​:分页查询 + 分批处理 → 超时,30分钟无响应
  • 最终目标​:30秒内完成30万数据导出

技术选型:为什么选择这种方案?

需求分析

  • 数据量:30万条记录
  • 字段数:15个字段/条
  • 文件格式:Excel
  • 响应时间:< 30秒
  • 内存占用:< 500MB

技术栈对比

方案 优点 缺点 适用场景
直接内存导出 实现简单 内存溢出风险 小数据量(<1万)
分页查询+流式导出 内存友好 数据库连接占用长 中等数据量
游标查询+流式处理 内存占用恒定 实现复杂 大数据量(推荐)

最终架构设计

用户请求 → 异步任务提交 → 游标查询 → 流式Excel生成 → 文件存储 → 下载通知

核心实现代码

1. 异步任务处理

@RestController
@Slf4j
public class DataExportController {
    
    @Autowired
    private ExportTaskService exportTaskService;
    
    @PostMapping("/export/carriers")
    public ResponseEntity<ExportResponse> exportCarriers(@RequestBody ExportRequest request) {
        log.info("接收到数据导出请求,参数: {}", request);
        
        String taskId = exportTaskService.submitExportTask(request);
        
        return ResponseEntity.ok(ExportResponse.success(taskId, "导出任务已提交,请稍后下载"));
    }
    
    @GetMapping("/export/download/{taskId}")
    public void downloadExportFile(@PathVariable String taskId, 
                                  HttpServletResponse response) {
        exportTaskService.downloadFile(taskId, response);
    }
    
    @GetMapping("/export/status/{taskId}")
    public ExportStatus getExportStatus(@PathVariable String taskId) {
        return exportTaskService.getTaskStatus(taskId);
    }
}

@Service
@Slf4j
public class ExportTaskService {
    
    @Autowired
    private ThreadPoolTaskExecutor exportTaskExecutor;
    
    @Autowired
    private ExportFileStorage fileStorage;
    
    private final Map<String, ExportStatus> taskStatusMap = new ConcurrentHashMap<>();
    
    public String submitExportTask(ExportRequest request) {
        String taskId = generateTaskId();
        
        // 初始状态
        ExportStatus status = new ExportStatus(taskId, "排队中", 0, null);
        taskStatusMap.put(taskId, status);
        
        // 异步执行导出任务
        exportTaskExecutor.execute(() -> {
            try {
                executeExportTask(taskId, request);
            } catch (Exception e) {
                log.error("导出任务执行失败: {}", taskId, e);
                updateTaskStatus(taskId, "失败", 100, e.getMessage());
            }
        });
        
        return taskId;
    }
    
    private void executeExportTask(String taskId, ExportRequest request) {
        updateTaskStatus(taskId, "数据查询中", 10, null);
        
        try (CarrierDataExporter exporter = new CarrierDataExporter(taskId, request)) {
            String filePath = exporter.export();
            updateTaskStatus(taskId, "完成", 100, filePath);
        }
    }
}

2. 游标查询 + 流式处理核心类

@Component
@Slf4j
public class CarrierDataExporter implements AutoCloseable {
    
    private final String taskId;
    private final ExportRequest request;
    private Connection connection;
    private PreparedStatement preparedStatement;
    private ResultSet resultSet;
    
    public CarrierDataExporter(String taskId, ExportRequest request) {
        this.taskId = taskId;
        this.request = request;
    }
    
    public String export() throws Exception {
        String filePath = generateFilePath();
        
        try (FileOutputStream fos = new FileOutputStream(filePath);
             SXSSFWorkbook workbook = new SXSSFWorkbook(100)) { // 保持100行在内存中
            
            // 创建Excel
            Sheet sheet = workbook.createSheet("载体数据");
            createHeaderRow(sheet);
            
            // 游标查询
            initializeCursorQuery();
            int rowNum = 1;
            int processedCount = 0;
            
            while (resultSet.next()) {
                if (rowNum % 1000 == 0) {
                    log.info("任务{}: 已处理{}条数据", taskId, rowNum);
                    updateProgress(taskId, rowNum);
                }
                
                // 流式写入Excel
                Row row = sheet.createRow(rowNum++);
                populateRowData(row, resultSet);
                processedCount++;
                
                // 每1000行刷新到磁盘
                if (processedCount % 1000 == 0) {
                    ((SXSSFSheet) sheet).flushRows(100); // 保留100行在内存中
                }
            }
            
            workbook.write(fos);
            log.info("导出任务完成,共处理{}条数据", processedCount);
            
        } finally {
            close();
        }
        
        return filePath;
    }
    
    private void initializeCursorQuery() throws SQLException {
        String sql = buildQuerySql();
        
        connection = DataSourceUtils.getConnection(dataSource);
        // 设置游标参数
        preparedStatement = connection.prepareStatement(sql, 
            ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY);
        preparedStatement.setFetchSize(1000); // 每次获取1000条
        
        setQueryParameters(preparedStatement, request);
        
        resultSet = preparedStatement.executeQuery();
    }
    
    private String buildQuerySql() {
        return """
            SELECT carrier_id, carrier_name, carrier_type, status, 
                   create_time, update_time, risk_level, total_orders,
                   active_days, region, industry, credit_score, 
                   contact_name, contact_phone, description
            FROM carrier_info 
            WHERE create_time >= ? AND create_time <= ?
            ORDER BY create_time DESC
            """;
    }
    
    private void populateRowData(Row row, ResultSet rs) throws SQLException {
        int colNum = 0;
        row.createCell(colNum++).setCellValue(rs.getString("carrier_id"));
        row.createCell(colNum++).setCellValue(rs.getString("carrier_name"));
        row.createCell(colNum++).setCellValue(rs.getString("carrier_type"));
        row.createCell(colNum++).setCellValue(rs.getString("status"));
        row.createCell(colNum++).setCellValue(formatDate(rs.getTimestamp("create_time")));
        // ... 其他字段
    }
    
    @Override
    public void close() {
        JdbcUtils.closeResultSet(resultSet);
        JdbcUtils.closeStatement(preparedStatement);
        DataSourceUtils.releaseConnection(connection, dataSource);
    }
}

3. 线程池和内存配置

@Configuration
@EnableAsync
public class AsyncExportConfig {
    
    @Bean("exportTaskExecutor")
    public ThreadPoolTaskExecutor exportTaskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(2);      // 核心线程数
        executor.setMaxPoolSize(5);       // 最大线程数
        executor.setQueueCapacity(10);     // 队列容量
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        executor.setThreadNamePrefix("export-task-");
        executor.setAwaitTerminationSeconds(60);
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.initialize();
        return executor;
    }
}

# application.yml 关键配置
spring:
  datasource:
    hikari:
      maximum-pool-size: 20
      connection-timeout: 30000
      idle-timeout: 600000
      max-lifetime: 1800000
  servlet:
    multipart:
      max-file-size: 500MB
      max-request-size: 500MB

server:
  max-http-header-size: 65536

# Excel导出配置
export:
  max-row-size: 1000000
  file-base-dir: /data/export/files
  file-keep-hours: 24

4. 数据库查询优化

-- 为导出查询创建合适的索引
CREATE INDEX idx_carrier_create_time ON carrier_info(create_time DESC);
CREATE INDEX idx_carrier_region_status ON carrier_info(region, status);

-- 查询优化:只选择需要的字段,避免SELECT *
EXPLAIN SELECT carrier_id, carrier_name, carrier_type, status, create_time 
FROM carrier_info 
WHERE create_time BETWEEN '2023-01-01' AND '2023-12-31'
ORDER BY create_time DESC;

5. 内存监控和优化

@Component
@Slf4j
public class ExportMemoryMonitor {
    
    @Scheduled(fixedRate = 5000) // 每5秒监控一次
    public void monitorMemoryUsage() {
        Runtime runtime = Runtime.getRuntime();
        long usedMemory = runtime.totalMemory() - runtime.freeMemory();
        long maxMemory = runtime.maxMemory();
        double usageRatio = (double) usedMemory / maxMemory;
        
        if (usageRatio > 0.8) {
            log.warn("内存使用率过高: {}/{} ({:.2%})", 
                formatMemory(usedMemory), formatMemory(maxMemory), usageRatio);
            // 触发GC(建议性)
            System.gc();
        }
        
        if (usageRatio > 0.9) {
            log.error("内存使用率超过90%,可能影响导出任务");
            // 可以在这里添加告警逻辑
        }
    }
    
    private String formatMemory(long bytes) {
        return String.format("%.2fMB", bytes / 1024.0 / 1024.0);
    }
}

性能测试结果

测试环境

  • 服务器:4核8G内存
  • 数据库:MySQL 8.0,单独服务器
  • 网络:千兆内网

性能数据

数据量 导出时间 内存占用 文件大小
1万条 2.3秒 150MB 3.2MB
10万条 8.5秒 180MB 31MB
30万条 22秒 220MB 93MB
50万条 38秒 250MB 155MB

与传统方案对比

方案 30万数据导出时间 内存占用 系统影响
传统分页查询 超时(>5分钟) 1.5GB+ 系统卡顿
直接内存加载 内存溢出 2GB+ 系统崩溃
本方案 22秒 220MB 几乎无影响

遇到的坑和解决方案

问题1:数据库连接超时

现象​:导出过程中数据库连接断开
解决​:调整数据库超时设置,使用游标查询

-- 调整会话超时时间
SET SESSION wait_timeout = 3600;
SET SESSION interactive_timeout = 3600;

问题2:Excel内存溢出

现象​:POI导出大文件时内存溢出
解决​:使用SXSSFWorkbook流式API

// 使用SXSSFWorkbook替代XSSFWorkbook
try (SXSSFWorkbook workbook = new SXSSFWorkbook(100)) { // 只保留100行在内存中
    // 流式处理
}

问题3:文件下载超时

现象​:大文件下载时网络超时
解决​:分块传输,支持断点续传

// 支持断点续传的文件下载
public void downloadFile(String taskId, HttpServletResponse response) {
    File file = fileStorage.getFile(taskId);
    
    long fileLength = file.length();
    long start = 0;
    long end = fileLength - 1;
    
    String rangeHeader = request.getHeader("Range");
    if (rangeHeader != null) {
        // 处理范围请求,支持断点续传
    }
    
    // 设置响应头
    response.setHeader("Accept-Ranges", "bytes");
    response.setContentType("application/vnd.openxmlformats-officedocument.spreadsheetml.sheet");
    response.setHeader("Content-Disposition", "attachment; filename=\"carriers.xlsx\"");
    
    // 流式传输文件内容
    try (FileInputStream fis = new FileInputStream(file);
         OutputStream os = response.getOutputStream()) {
        
        byte[] buffer = new byte[8192];
        int bytesRead;
        while ((bytesRead = fis.read(buffer)) != -1) {
            os.write(buffer, 0, bytesRead);
        }
    }
}

最佳实践总结

1. 架构设计原则

  • 异步处理​:避免长时间阻塞HTTP请求
  • 流式处理​:内存占用恒定,不受数据量影响
  • 进度反馈​:让用户了解任务状态
  • 错误处理​:完善的异常处理和重试机制

2. 性能优化要点

  • 游标查询​:避免一次性加载所有数据到内存
  • 分批处理​:控制每次处理的数据量
  • 内存管理​:及时释放资源,监控内存使用
  • 索引优化​:确保查询性能

3. 用户体验优化

  • 进度展示​:实时显示导出进度
  • 错误提示​:明确的错误信息和解决方案
  • 文件管理​:自动清理过期文件
  • 断点续传​:支持大文件断点下载

未来优化方向

  1. 分布式导出​:对于亿级数据,考虑分布式处理
  2. 压缩优化​:支持更高效的文件压缩格式
  3. 缓存预热​:预生成常用查询的导出文件
  4. 智能分片​:根据数据特征自动优化查询策略

结语

通过这次30万数据导出方案的实践,我们成功将导出时间从超时优化到22秒,内存占用控制在220MB以内。关键技术点包括游标查询、流式Excel处理、异步任务管理等。

大数据导出不仅是技术挑战,更是对系统架构设计的考验。合理的架构设计能够在保证性能的同时,提供良好的用户体验。

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐