You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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};
        }
    }
}

关键修改点说明

  1. 取消循环内的阻塞调用:先批量提交所有任务到线程池,收集Future列表后再统一处理结果,让线程池的多个线程能同时执行任务,真正实现并发。
  2. 线程安全的资源隔离:
    • 每个CallWebService任务自行创建和关闭SOAPConnection,避免多线程共享导致的线程安全问题
    • 数据库更新使用独立的Connection(生产环境更推荐用数据库连接池,比如HikariCP,减少连接创建销毁的开销)
    • 不再在任务间共享PreparedStatement,更新操作复用同一个PreparedStatement提升效率
  3. 资源泄漏防护:使用try-with-resources语法自动管理数据库连接和语句,确保资源能被正确关闭
  4. 线程池优雅关闭:增加shutdownNow()兜底逻辑,确保程序在超时情况下能正常退出

内容的提问来源于stack exchange,提问作者vamsi krishna

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 08:45:35