Java Scheduled Executor线程进入等待状态问题求助
问题分析与解决方法
首先咱们拆解下你遇到的核心问题:双线程调度池下程序执行几次读写后陷入无限等待,但单线程调度池能正常运行。结合你的代码和线程栈信息,我来逐一分析原因和对应的解决办法:
核心原因1:ScheduledExecutorService的使用方式错误
你用ScheduledExecutorService.schedule()提交线程任务,但这个方法是用来延迟/定时执行一次性任务的,而你只是需要立即启动两个并发线程,属于典型的误用:
- 单线程调度池(
newSingleThreadScheduledExecutor())会串行执行任务:先跑完Writer的2000次请求,再启动Reader读响应,这时候服务器已经处理完所有请求,Reader能顺利读到所有结果。 - 双线程调度池(
newScheduledThreadPool(2))会同时启动Writer和Reader:Writer快速发送请求,Reader边读边等,但Writer跑完任务后,对应的线程会回到调度池的等待队列(这就是你线程栈里pool-1-thread-2处于WAITING状态的原因,本身是正常的),而Reader会一直阻塞在reader.readLine()——如果服务器没主动关闭连接,它就会无限等待下去。
核心原因2:Reader线程未处理连接关闭的情况
当服务器处理完所有请求并关闭Socket连接时,reader.readLine()会返回null,但你的Reader线程只是打印提示后继续循环,不会退出,导致程序看起来一直在“等待”。
次要问题:Redis键过期时间过短
你给Redis的请求ID设置了10秒过期,如果服务器处理请求速度较慢,Reader收到响应时对应的Key已经过期,会导致日志丢失,但不会直接导致程序卡住。
解决方法
1. 替换ScheduledExecutorService为普通线程池或直接启动线程
你不需要调度任务,只需要并发执行两个线程,推荐两种方案:
方案A:直接启动Thread(最简单)
把主线程里的调度池代码替换成直接启动线程:
public static void main(String[] args) throws IOException { Socket socket = new Socket(HOSTNAME, PORT); PrintWriter writer = new PrintWriter(socket.getOutputStream()); BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream())); Jedis jedis = new Jedis("localhost"); WithdrawRequestReader readerThread = new WithdrawRequestReader(reader, jedis); WithdrawRequestWriter writerThread = new WithdrawRequestWriter(writer, jedis); // 直接启动两个线程 writerThread.start(); readerThread.start(); }
方案B:用普通线程池(更规范)
如果想用线程池管理,用ExecutorService而非ScheduledExecutorService,用execute()提交任务:
public static void main(String[] args) throws IOException { Socket socket = new Socket(HOSTNAME, PORT); PrintWriter writer = new PrintWriter(socket.getOutputStream()); BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream())); Jedis jedis = new Jedis("localhost"); // 创建固定大小的线程池 ExecutorService executor = Executors.newFixedThreadPool(2); WithdrawRequestReader readerThread = new WithdrawRequestReader(reader, jedis); WithdrawRequestWriter writerThread = new WithdrawRequestWriter(writer, jedis); // 提交任务立即执行 executor.execute(writerThread); executor.execute(readerThread); // 可选:如果需要等待所有任务完成后关闭程序,可添加 // executor.shutdown(); // executor.awaitTermination(5, TimeUnit.MINUTES); }
2. 修改Reader线程逻辑,处理连接关闭
当服务器关闭连接时,Reader线程应该退出循环,避免无限等待:
@Override public void run() { String response; try { while ((response = reader.readLine()) != null) { String[] tokenizedResponse = response.split(" "); System.out.println("Response is: " + response); String value = jedis.get(tokenizedResponse[0]); if (value != null) { System.out.println("Request Token: " + tokenizedResponse[0] + " | Status: " + tokenizedResponse[1] + " | Account: " + tokenizedResponse[2] + " | Amount: " + tokenizedResponse[3] + " | Remaining Balance: " + tokenizedResponse[4] + " | ResponeStr: " + tokenizedResponse[5]); } } // readLine()返回null说明服务器关闭了连接 System.out.println("Server closed connection. Exiting reader thread."); } catch (IOException ex) { ex.printStackTrace(); } finally { // 关闭Redis连接 if (jedis != null) { jedis.close(); } } }
3. 优化Writer线程,发送完请求后关闭输出流
在Writer线程的循环结束后,关闭PrintWriter,告诉服务器已无更多请求,这样服务器可以处理完剩余请求后主动关闭连接:
@Override public void run() { try { int i = 0; for (int accountNo = START_ACCOUNT_NO; accountNo <= END_ACCOUNT_NO; accountNo++) { String uniqueId = UUID.randomUUID().toString(); // 把过期时间调长到60秒,避免响应回来时Key已过期 if (jedis.set(uniqueId, String.valueOf(accountNo), "NX", "EX", 60) != null) { writer.write(uniqueId + " " + accountNo + " " + AMOUNT_TO_WITHDRAW + "\n"); writer.flush(); } else { System.err.println("Setting of UUID in redis failed: " + uniqueId); } } // 发送完所有请求,关闭输出流 writer.close(); System.out.println("All requests sent. Closed output stream."); } finally { // 关闭Redis连接 if (jedis != null) { jedis.close(); } } }
4. 调整Redis键过期时间
把Redis的Key过期时间从10秒调整到60秒(或更长,根据服务器处理速度),避免响应返回时Key已过期导致日志丢失。
内容的提问来源于stack exchange,提问作者ks2bmallik
相关产品推荐
相关产品推荐

