Java 7多线程模块引发数据库死锁问题及优化咨询
问题根源分析
你遇到的不同线程获取相同数据的问题,核心原因是所有异步任务共享了同一个CallableModel实例!
看你的代码逻辑:创建ThreadConcurrentWoker时只实例化了一个CallableModel匿名类对象,循环提交任务时,所有线程都调用这个共享实例的setElement方法。由于setElement没有任何线程安全保护,再加上线程调度的不确定性,很容易出现以下场景:
- 线程A调用
setElement(elemA),还没来得及执行call()就被CPU暂停 - 线程B抢占CPU,调用
setElement(elemB)覆盖了共享实例里的元素值 - 线程A恢复执行时,
call()拿到的是elemB而非自己的elemA
最终导致多个线程重复处理同一个数据,触发数据库重复操作,进而引发死锁。
改进方案
解决问题的核心思路是:让每个异步任务拥有独立的状态实例,彻底避免共享变量带来的线程安全问题。下面提供两种可行的改进方式:
方式一:修改任务提交逻辑,为每个元素创建独立的CallableModel
在遍历目标列表时,不为所有任务复用同一个CallableModel,而是为每个元素创建专属的实例。这样每个线程操作的都是自己的对象,不会出现状态覆盖。
修改后的ThreadConcurrentWoker.concurrentExcute()方法代码:
@Override public List<R> concurrentExcute() throws Exception { // 线程池大小建议设上限,避免列表过大时创建大量线程耗尽资源 int poolSize = Math.min(super.targetList.size(), Runtime.getRuntime().availableProcessors() * 2); ExecutorService executor = Executors.newFixedThreadPool(poolSize); CompletionService<R> completionService = new ExecutorCompletionService<>(executor); for (final E element : super.targetList) { // 为每个元素创建独立的CallableModel实例 completionService.submit(new CallableModel<E, R>() { @Override public R call() throws Exception { // 直接使用当前循环的element,无需通过setElement传递 // 这里编写你的业务逻辑 ResultBean result = new ResultBean(); // do something with element return result; } }); } int finishs = 0; boolean errors = false; while (finishs < super.targetList.size() && !errors) { Future<R> resultFuture = completionService.take(); try { super.results.add(resultFuture.get()); } catch (ExecutionException e) { errors = true; // 建议记录完整异常栈,方便排查问题 log.error("任务执行失败", e); } finally { finishs++; } } // 关闭线程池,释放资源 executor.shutdown(); return super.results; }
方式二:重构CallableModel,改用函数式接口传递业务逻辑
我们可以去掉CallableModel的成员变量element,用函数式接口直接传递业务逻辑,从根源上消除共享状态。
首先定义一个函数式接口替代原有的CallableModel:
@FunctionalInterface public interface TaskHandler<E, R> { R handle(E element) throws Exception; }
然后修改ThreadConcurrentWoker的构造和执行逻辑:
public class ThreadConcurrentWoker<E, R> extends ThreadConcurrent<E, R> { private final TaskHandler<E, R> taskHandler; public ThreadConcurrentWoker(List<E> target, TaskHandler<E, R> taskHandler) { super.targetList = target; this.taskHandler = taskHandler; super.results = new Vector<>(); } @Override public List<R> concurrentExcute() throws Exception { int poolSize = Math.min(super.targetList.size(), Runtime.getRuntime().availableProcessors() * 2); ExecutorService executor = Executors.newFixedThreadPool(poolSize); CompletionService<R> completionService = new ExecutorCompletionService<>(executor); for (final E element : super.targetList) { // 每个任务直接持有自己的element,通过函数式接口执行逻辑 completionService.submit(() -> taskHandler.handle(element)); } int finishs = 0; boolean errors = false; while (finishs < super.targetList.size() && !errors) { Future<R> resultFuture = completionService.take(); try { super.results.add(resultFuture.get()); } catch (ExecutionException e) { errors = true; log.error("任务执行失败", e); } finally { finishs++; } } executor.shutdown(); return super.results; } }
使用时的代码会更简洁:
ThreadConcurrentWoker<FlowPendingCheckedBean, ResultBean> tCUtil = new ThreadConcurrentWoker<>(test, element -> { // 直接处理当前element,编写业务逻辑 ResultBean result = new ResultBean(); // do something with element return result; }); try { resultBeans = tCUtil.concurrentExcute(); } catch (Exception e1) { log.error("多线程执行失败", e1); }
额外优化建议
- 线程池大小控制:不要直接用
targetList.size()作为线程池大小,当列表元素过多时会创建大量线程,导致系统资源耗尽。建议根据CPU核心数设置合理上限。 - 异常处理策略:原代码中遇到异常就终止后续任务,你可以根据业务需求调整,比如记录异常后继续处理其他任务。
- 资源释放:务必调用
executor.shutdown()关闭线程池,避免线程泄漏。
内容的提问来源于stack exchange,提问作者J. Hu
相关产品推荐
相关产品推荐

