使用ExecutorService实现多线程却串行执行问题排查求助
问题根源与解决建议
你的代码之所以看起来是串行执行,核心问题出在循环内提交任务后立即调用future.get()——这个方法会阻塞当前主线程,必须等到当前任务执行完毕才会继续下一次循环,相当于把多线程变成了排队执行的串行任务。
除此之外,代码里还有几个线程安全的隐患必须修复,否则即使改成并发也会出现奇怪的问题:
- JDBC的
Connection和PreparedStatement不是线程安全的,不能在多个Callable任务之间共享 SOAPConnection同样不是线程安全的,每个任务应该独立创建和关闭
修复后的代码示例
import java.io.BufferedReader; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; import java.io.Reader; import java.sql.Clob; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Statement; import java.util.ArrayList; import java.util.Date; import java.util.List; import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.apache.commons.io.input.ReaderInputStream; import javax.xml.soap.*; public class SoapClientSVAPI2 { // 提取配置常量,方便维护 private static final String DB_URL = "jdbc:oracle:thin:@XXX:1521:XXX"; private static final String DB_USER = "XXX"; private static final String DB_PWD = "XX"; private static final String SOAP_ENDPOINT_URL = "XXX"; public static void main(String args[]) { Connection conn = null; Statement s = null; ExecutorService executor = Executors.newFixedThreadPool(10); // 用来保存所有任务的Future对象,后续统一处理结果 List<Future<String[]>> futureList = new ArrayList<>(); try { Class.forName("oracle.jdbc.driver.OracleDriver"); conn = DriverManager.getConnection(DB_URL, DB_USER, DB_PWD); s = conn.createStatement(); String migTable = "SELECT req_id,payload,attr_9 FROM tableX WHERE ATTR_9 = 'READY' ORDER BY dbms_random.value"; ResultSet rs = s.executeQuery(migTable); // 第一步:批量提交所有任务到线程池,不立即阻塞等待结果 while (rs != null && rs.next()) { Clob payload = rs.getClob(2); String reqid = rs.getString(1); if (reqid != null && !reqid.isEmpty()) { // 每个任务独立创建Callable,不共享外部资源 CallWebService cweb = new CallWebService(reqid, payload); Future<String[]> future = executor.submit(cweb); futureList.add(future); } } // 第二步:所有任务提交完成后,统一处理结果并更新数据库 // 使用try-with-resources自动管理更新用的数据库资源 try (Connection updateConn = DriverManager.getConnection(DB_URL, DB_USER, DB_PWD); PreparedStatement updateStmt = updateConn.prepareStatement( "UPDATE tableX SET RESPONSE = ? , attr_9 = ? WHERE REQ_ID = ?")) { for (Future<String[]> future : futureList) { String reqidout = ""; String response = ""; String status = "FAILURE"; try { String[] result = future.get(10, TimeUnit.SECONDS); reqidout = result[0]; response = result[1]; if (response != null && response.indexOf("ReturnCode") > 0) { status = "SUCCESS"; } // 执行数据库更新 updateStmt.setString(1, response); updateStmt.setString(2, status); updateStmt.setString(3, reqidout); updateStmt.executeUpdate(); } catch (InterruptedException e) { e.printStackTrace(); response = "InterruptedException"; } catch (ExecutionException e) { e.printStackTrace(); response = "ExecutionException"; } catch (TimeoutException e) { e.printStackTrace(); response = "TimeoutException"; future.cancel(true); } } } } catch (Exception e) { System.out.println("\nError occurred while initializing resources!\n"); e.printStackTrace(); } finally { System.out.println("Time end" + new Date()); // 优雅关闭线程池 executor.shutdown(); try { if (!executor.awaitTermination(1, TimeUnit.HOURS)) { executor.shutdownNow(); } } catch (InterruptedException e1) { executor.shutdownNow(); e1.printStackTrace(); } // 关闭初始的数据库连接 try { if (s != null) s.close(); if (conn != null) conn.close(); } catch (SQLException e) { e.printStackTrace(); } } } // 独立的Callable任务类,每个任务自行管理专属资源 static class CallWebService implements Callable<String[]> { private final String reqid; private final Clob payload; public CallWebService(String reqid, Clob payload) { this.reqid = reqid; this.payload = payload; } @Override public String[] call() throws Exception { ByteArrayOutputStream baos = null; SOAPConnection soapConn = null; String svresponse = ""; try { String threadName = Thread.currentThread().getName(); System.out.println("Reqid: " + reqid + " Thread name:" + threadName); // 每个任务独立创建SOAP连接 SOAPConnectionFactory soapConnectionFactory = SOAPConnectionFactory.newInstance(); soapConn = soapConnectionFactory.createConnection(); Reader ir = new BufferedReader(payload.getCharacterStream()); InputStream is = new ReaderInputStream(ir, "UTF-8"); SOAPMessage inpRequest = MessageFactory.newInstance().createMessage(null, is); SOAPMessage soapResponse = soapConn.call(inpRequest, SOAP_ENDPOINT_URL); if (soapResponse != null) { baos = new ByteArrayOutputStream(); soapResponse.writeTo(baos); svresponse = baos.toString(); } } catch (Exception e) { e.printStackTrace(); svresponse = "ERROR: " + e.getMessage(); } finally { // 关闭当前任务的所有资源 if (baos != null) { try { baos.close(); } catch (IOException ioe) { ioe.printStackTrace(); } } if (soapConn != null) { try { soapConn.close(); } catch (SOAPException e) { e.printStackTrace(); } } } return new String[]{reqid, svresponse}; } } }
关键修改点说明
- 取消循环内的阻塞调用:先批量提交所有任务到线程池,收集Future列表后再统一处理结果,让线程池的多个线程能同时执行任务,真正实现并发。
- 线程安全的资源隔离:
- 每个
CallWebService任务自行创建和关闭SOAPConnection,避免多线程共享导致的线程安全问题 - 数据库更新使用独立的Connection(生产环境更推荐用数据库连接池,比如HikariCP,减少连接创建销毁的开销)
- 不再在任务间共享
PreparedStatement,更新操作复用同一个PreparedStatement提升效率
- 每个
- 资源泄漏防护:使用try-with-resources语法自动管理数据库连接和语句,确保资源能被正确关闭
- 线程池优雅关闭:增加
shutdownNow()兜底逻辑,确保程序在超时情况下能正常退出
内容的提问来源于stack exchange,提问作者vamsi krishna
相关产品推荐
相关产品推荐

