多线程场景下使用JSch执行命令部分返回结果不完整问题咨询
问题原因
- 资源关闭时序错误,未等待远程命令执行完成就终止读取:JSch的
ChannelExec的输入流会在命令执行过程中持续返回数据,你当前仅依赖readLine()!=null判断读取结束,多线程场景下CPU调度切换会导致还未读取完所有输出就跳出循环,提前关闭输入流和通道。单线程场景下命令执行和读取的时序刚好匹配,所以能拿到完整结果。 - 数据存储逻辑缺失:你在循环中读取到
lineData并做字符串替换后,没有将处理后的内容添加到预先定义的dataList集合中,读取到的内容会直接被丢弃。 - 错误流未收集:部分命令的输出可能会写入标准错误流,你直接将错误流重定向到
System.err没有做收集,也会导致部分结果丢失。 - 线程不安全类混用:代码中用到的
sd是SimpleDateFormat实例,该类本身线程不安全,多线程同时调用format方法会产生不可预期的异常,可能干扰正常的执行流程。
解决方案
按照以下要求修改代码即可解决输出缺失问题:
- 调整读取逻辑,循环判断通道是否关闭,同时持续读取输入流中所有可用数据,等待命令执行完成后再关闭资源
- 修复数据存储逻辑,将处理后的行数据存入集合
- 单独收集标准错误流的输出
- 替换线程不安全的日期格式化类
修改后的参考代码如下:
// 替换SimpleDateFormat为线程安全的DateTimeFormatter,pattern替换为你实际使用的日期格式 DateTimeFormatter dtf = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); ExecutorService executor = Executors.newFixedThreadPool(3); List<CompletableFuture<Void>> cfs = new ArrayList<>(); AppConfig.nimblesMap.forEach((k, v)->{ CompletableFuture<Void> cf = CompletableFuture.runAsync(()->{ System.out.println("Schedule start,Machine Name:"+k+"; Date Time:"+ dtf.format(LocalDateTime.now())); String machineName = v.getMachineName(); Session session = null; ChannelExec channelExec = null; BufferedReader inputBr = null; BufferedReader errorBr = null; try { // 补全主机、端口参数 session = new JSch().getSession(v.getUserName(), v.getHost(), v.getPort()); session.setPassword(v.getPassWord()); Properties config = new Properties(); config.put("StrictHostKeyChecking", "no"); session.setConfig(config); session.connect(10000); // 增加连接超时时间,避免无限等待 channelExec = (ChannelExec) session.openChannel("exec"); channelExec.setCommand("你要执行的具体命令"); channelExec.setInputStream(null); // 单独获取错误流,不直接重定向到System.err InputStream errIn = channelExec.getErrStream(); InputStream in = channelExec.getInputStream(); channelExec.connect(5000); // 通道连接超时 inputBr = new BufferedReader(new InputStreamReader(in)); errorBr = new BufferedReader(new InputStreamReader(errIn)); ArrayList<String> dataList = new ArrayList<>(); ArrayList<String> errorList = new ArrayList<>(); String lineData; // 等待通道关闭,持续读取数据 while (!channelExec.isClosed()) { // 读取标准输出 while ((lineData = inputBr.readLine()) != null) { lineData = lineData.replaceAll("\\s+", ","); dataList.add(lineData); // 补全数据存储逻辑 } // 读取错误输出 while ((lineData = errorBr.readLine()) != null) { errorList.add(lineData); } // 避免CPU空转,增加短暂休眠 Thread.sleep(100); } // 通道关闭后再读一次剩余数据 while ((lineData = inputBr.readLine()) != null) { lineData = lineData.replaceAll("\\s+", ","); dataList.add(lineData); } while ((lineData = errorBr.readLine()) != null) { errorList.add(lineData); } // 此处补充你自己的业务逻辑,处理收集到的dataList和errorList } catch (Exception e) { e.printStackTrace(); } finally { // 资源在finally块关闭,避免泄漏 if (inputBr != null) { try {inputBr.close();} catch (Exception ignored) {} } if (errorBr != null) { try {errorBr.close();} catch (Exception ignored) {} } if (channelExec != null && channelExec.isConnected()) { channelExec.disconnect(); } if (session != null && session.isConnected()) { session.disconnect(); } } },executor); cfs.add(cf); }); // 统一等待所有任务执行完成 CompletableFuture.allOf(cfs.toArray(new CompletableFuture[0])).join(); executor.shutdown();
内容的提问来源于stack exchange,提问作者Leo Liu
相关产品推荐
相关产品推荐

