SpringBoot多线程批量插入MySQL数据未持久化问题排查
尝试通过多线程向MySQL插入50000条数据,日志中未发现任何异常或错误,但数据仍未持久化到数据库,以下是待排查的代码:
@Retryable( value = { Exception.class }, maxAttempts = 3, backoff = @Backoff(delay = 5000)) @Transactional public void insertBatch(List<Customer> batch) { int numThreads = 4; // Set the number of threads int batchSize = 1000; // Set the batch size List<List<Customer>> partitions = new ArrayList<>(); for (int i = 0; i < batch.size(); i += batchSize) { int end = Math.min(i + batchSize, batch.size()); partitions.add(batch.subList(i, end)); } ExecutorService executor = Executors.newFixedThreadPool(numThreads); List<Future<?>> futures = new ArrayList<>(); for (List<Customer> partition : partitions) { futures.add(executor.submit(() -> { try { customerRepository.saveAll(partition); log.info("Customer Saved "+ partition.size()); } catch (Exception e) { log.error("Failed to insert partition: {}", partition, e); throw e; } })); } for (Future<?> future : futures) { try { future.get(); } catch (Exception e) { log.error("Partition insertion failed", e); } } executor.shutdown(); } @Recover public void handleInsertBatchException(Exception e, List<Customer> batch) { log.error("Insertion failed for batch: {}", batch, e); // Handle the exception here, for example, you can retry with a smaller batch size }
输出日志:
日志中显示多条Customer Saved 1000的成功记录,最终出现Insertion failed for batch: [...]的错误日志。
问题排查与修复方案
1. 事务与多线程的冲突
主线程添加的@Transactional注解无法作用于子线程:Spring的事务基于线程上下文,子线程不会继承主线程的事务上下文。主线程的事务仅覆盖自身的线程调度逻辑,子线程中customerRepository.saveAll()会开启独立事务,但主线程的事务状态完全不影响子线程。这种情况下,主线程的事务注解完全多余,甚至可能导致逻辑混淆。
修复:移除insertBatch方法上的@Transactional注解。
2. 异常处理导致重试机制失效
遍历Future时,捕获异常仅打印日志但未重新抛出,导致@Retryable注解无法感知子线程的异常,不会触发重试逻辑。同时主线程无法传递子线程的异常,可能掩盖真实失败原因。
修复:捕获异常后重新抛出,确保主线程能感知异常并触发重试:
for (Future<?> future : futures) { try { future.get(); } catch (Exception e) { log.error("Partition insertion failed", e); throw new RuntimeException(e); // 重新抛出异常,触发重试 } }
3. subList的视图风险
batch.subList(i, end)返回的是原列表的视图而非独立新列表,若原列表被其他操作修改,会直接影响子列表内容,可能导致JPA保存时出现数据不一致。
修复:创建新ArrayList存储分片数据:
partitions.add(new ArrayList<>(batch.subList(i, end)));
4. 线程池关闭的不严谨性
仅调用executor.shutdown()无法确保所有线程任务完成后再结束,需添加等待逻辑避免主线程提前退出导致子线程任务中断。
修复:添加awaitTermination等待线程池任务完成:
executor.shutdown(); try { if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { executor.shutdownNow(); if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { log.error("Executor service did not terminate"); } } } catch (InterruptedException ie) { executor.shutdownNow(); Thread.currentThread().interrupt(); }
修改后的完整代码示例
@Retryable( value = { Exception.class }, maxAttempts = 3, backoff = @Backoff(delay = 5000)) public void insertBatch(List<Customer> batch) { int numThreads = 4; int batchSize = 1000; List<List<Customer>> partitions = new ArrayList<>(); for (int i = 0; i < batch.size(); i += batchSize) { int end = Math.min(i + batchSize, batch.size()); partitions.add(new ArrayList<>(batch.subList(i, end))); } ExecutorService executor = Executors.newFixedThreadPool(numThreads); List<Future<?>> futures = new ArrayList<>(); for (List<Customer> partition : partitions) { futures.add(executor.submit(() -> { try { customerRepository.saveAll(partition); log.info("Customer Saved "+ partition.size()); } catch (Exception e) { log.error("Failed to insert partition: {}", partition, e); throw e; } })); } for (Future<?> future : futures) { try { future.get(); } catch (Exception e) { log.error("Partition insertion failed", e); throw new RuntimeException(e); } } executor.shutdown(); try { if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { executor.shutdownNow(); if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { log.error("Executor service did not terminate"); } } } catch (InterruptedException ie) { executor.shutdownNow(); Thread.currentThread().interrupt(); } } @Recover public void handleInsertBatchException(Exception e, List<Customer> batch) { log.error("Insertion failed for batch: {}", batch, e); // 可在此处添加自定义降级逻辑,例如改用单线程批量插入 }
内容的提问来源于stack exchange,提问作者newcoder

