diff --git a/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/service/impl/CSCollectServiceImpl.java b/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/service/impl/CSCollectServiceImpl.java index b7e339ed..8e7f805d 100644 --- a/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/service/impl/CSCollectServiceImpl.java +++ b/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/service/impl/CSCollectServiceImpl.java @@ -29,6 +29,7 @@ import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; @Service @Slf4j @@ -64,14 +65,19 @@ public class CSCollectServiceImpl implements CSCollectService { */ @Override public JSONArray getHpcAppCluster(String appCode) throws Exception { - String url = csUrl + "/pcm/v1/hpc/getHpcAppCluster"; - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - String result = HttpUtils.sendGetWithToken(url, "app=" + appCode, csAuthInfo.getToken()); - JSONObject resData = returnResData(result); - JSONArray list = (JSONArray) getByPath(resData, "list"); + try { + String url = csUrl + "/pcm/v1/hpc/getHpcAppCluster"; + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + String result = HttpUtils.sendGetWithToken(url, "app=" + appCode, csAuthInfo.getToken()); + JSONObject resData = returnResData(result); + JSONArray list = (JSONArray) getByPath(resData, "list"); // JSONObject param = JSONUtil.parseObj("{\"id\":\"11111" + "\",\"name\":\"自由调度\",\"nickname\":\"自由调度\",\"region\":\"\"}"); // return list.put(param); - return list; + return list; + } catch (Exception e) { + log.error("getHpcAppCluster error", e); + throw new Exception("getHpcAppCluster error:" + e); + } } @@ -85,48 +91,57 @@ public class CSCollectServiceImpl implements CSCollectService { */ @Override public Integer uploadFile(String appCode, MultipartFile file) throws Exception { - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - String token = csAuthInfo.getToken(); - Integer userId = csAuthInfo.getUserId(); - //1.获取packageId - Integer packageId = queryPackageId(appCode, userId, token); - //2.文件鉴权 - String url = presign(userId, token, packageId, file.getOriginalFilename()); - HashMap headers = MapUtil.newHashMap(); - headers.put("authorization", "Bearer " + csAuthInfo.getToken()); - headers.put("content-type", "application/octet-stream"); - String res = HttpRequest.post(url) - .addHeaders(headers) - .body(file.getBytes()) - .contentType("application/json") - .execute() - .body(); - return packageId; + try { + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + String token = csAuthInfo.getToken(); + Integer userId = csAuthInfo.getUserId(); + //1.获取packageId + Integer packageId = queryPackageId(appCode, userId, token); + //2.文件鉴权 + String url = presign(userId, token, packageId, file.getOriginalFilename()); + HashMap headers = MapUtil.newHashMap(); + headers.put("authorization", "Bearer " + csAuthInfo.getToken()); + headers.put("content-type", "application/octet-stream"); + String res = HttpRequest.post(url) + .addHeaders(headers) + .body(file.getBytes()) + .contentType("application/json") + .execute() + .body(); + return packageId; + } catch (Exception e) { + log.error("uploadFile error", e); + throw new Exception("uploadFile error:" + e); + } } /** * 任务提交执行 * - * @throws Exception - * example + * @throws Exception example */ @Override public String submit(TrainingTask trainingTask) throws Exception { - JSONObject taskInfo = JSONUtil.parseObj(trainingTask.getTaskInfo()); - JSONObject taskParam = JSONUtil.parseObj(trainingTask.getTaskParam()); - String clusterId = (String) taskInfo.get("type"); - String ntasks = (String) taskInfo.get("ntasks"); - String nodes = (String) taskInfo.get("nodes"); - String jobName = taskInfo.get("name").toString(); + try { + JSONObject taskInfo = JSONUtil.parseObj(trainingTask.getTaskInfo()); + JSONObject taskParam = JSONUtil.parseObj(trainingTask.getTaskParam()); + String clusterId = (String) taskInfo.get("type"); + String ntasks = (String) taskInfo.get("ntasks"); + String nodes = (String) taskInfo.get("nodes"); + String jobName = taskInfo.get("name").toString(); - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - List> hpcBindingFiles = buildHpcBindingFiles(taskParam, taskInfo); - JSONObject params = convertToJobRequest(trainingTask.getTaskName(),clusterId, jobName, hpcBindingFiles, csAuthInfo.getUserId(), ntasks, nodes); - String url = csUrl + "/jsm/v2/jobs/submit"; - String res = HttpUtils.sendBodyPostWithToken(url, params, csAuthInfo.getToken()); - JSONObject resData = returnResDataWithOK(res); - JSONObject taskIdObj = JSONUtil.parseObj(getByPath(resData, "message")); - return getByPath(taskIdObj, "jobInfo.taskId").toString(); + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + List> hpcBindingFiles = buildHpcBindingFiles(trainingTask.getTaskName(), taskParam, taskInfo); + JSONObject params = convertToJobRequest(trainingTask.getTaskName(), clusterId, jobName, hpcBindingFiles, csAuthInfo.getUserId(), ntasks, nodes); + String url = csUrl + "/jsm/v2/jobs/submit"; + String res = HttpUtils.sendBodyPostWithToken(url, params, csAuthInfo.getToken()); + JSONObject resData = returnResDataWithOK(res); + JSONObject taskIdObj = JSONUtil.parseObj(getByPath(resData, "message")); + return getByPath(taskIdObj, "jobInfo.taskId").toString(); + } catch (Exception e) { + log.error("submit error", e); + throw new Exception("submit error:" + e); + } } @@ -138,20 +153,25 @@ public class CSCollectServiceImpl implements CSCollectService { */ @Override public JSONObject queryTaskDetails(String taskId) throws Exception { - String url = csUrl + "/pcm/v1/core/task/details"; + try { + String url = csUrl + "/pcm/v1/core/task/details"; // String url = csUrl + "/jsm/v2/jobs/details"; - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - Integer userId = csAuthInfo.getUserId(); - String token = csAuthInfo.getToken(); - String res = HttpUtils.sendGetWithToken(url, "id=" + taskId, token); + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + Integer userId = csAuthInfo.getUserId(); + String token = csAuthInfo.getToken(); + String res = HttpUtils.sendGetWithToken(url, "id=" + taskId, token); // JSONObject resData = returnResDataWithOK(res); - JSONObject resData = returnResData(res); - //1.拼凑回显字段 - //2.检查任务运行状态 - String status = checkRunStatus(getByPath(resData, "subTaskInfos[0].status").toString()); - //3.若任务执行成功对task_result填入值 - resData.set("status", status); - return resData; + JSONObject resData = returnResData(res); + //1.拼凑回显字段 + //2.检查任务运行状态 + String status = checkRunStatus(getByPath(resData, "subTaskInfos[0].status").toString()); + //3.若任务执行成功对task_result填入值 + resData.set("status", status); + return resData; + } catch (Exception e) { + log.error("queryTaskDetails error", e); + throw new Exception("queryTaskDetails error:" + e); + } } /** @@ -163,18 +183,112 @@ public class CSCollectServiceImpl implements CSCollectService { */ @Override public String getJobLogs(Long taskId) throws Exception { - String url = csUrl + "/pcm/v1/hpc/jobLogs/" + taskId; - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - String res = HttpUtils.sendGetWithToken(url, null, csAuthInfo.getToken()); - JSONObject resData = returnResData(res); - return resData.get("outLogs").toString(); + try { + String url = csUrl + "/pcm/v1/hpc/jobLogs/" + taskId; + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + String res = HttpUtils.sendGetWithToken(url, null, csAuthInfo.getToken()); + JSONObject resData = returnResData(res); + return resData.get("outLogs").toString(); + } catch (Exception e) { + log.error("getJobLogs error", e); + throw new Exception("getJobLogs error:" + e); + } + } + + + /** + * 任务执行结果查询 + * + * @param taskId + * @throws Exception + */ + @Override + public JSONArray queryTaskResult(String taskId) throws Exception { + try { + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + String url = csUrl + "/jsm/v2/jobs/results"; + JSONObject params = new JSONObject(); + params.set("taskID", taskId); + params.set("type", "HPCSlurm"); + params.set("userID", csAuthInfo.getUserId()); + String res = HttpUtils.sendBodyPostWithToken(url, params, csAuthInfo.getToken()); + JSONObject resData = returnResDataWithOK(res); + return JSONUtil.parseArray(getByPath(resData, "pcmJobData.resultFiles[0].objects")); + } catch (Exception e) { + log.error("queryTaskResult error", e); + throw new Exception("queryTaskResult error:" + e); + } + } + + + /** + * 单个文件下载 + * + * @param objectId + * @param savePath + * @param fileName + * @throws Exception + */ + @Override + public JSONObject singleFileDownload(Integer objectId, String savePath, String fileName) throws Exception { + try { + CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); + String url = csUrl + "/jsm/v2/storage/presign"; + String paramStr = "{\"userID\":137,\"info\":{\"type\":\"download\",\"params\":{\"objectID\":70805}}}"; + JSONObject params = JSONUtil.parseObj(paramStr); + params.set("userID", csAuthInfo.getUserId()); + params.putByPath("info.params.objectID", objectId); + //获取下载地址 + String res = HttpUtils.sendBodyPostWithToken(url, params, csAuthInfo.getToken()); + JSONObject resData = returnResDataWithOK(res); + String presignUrl = resData.get("presignUrl").toString(); + //下载文件 +// CloseableHttpResponse response = HttpUtils.sendGetReturnResponse(presignUrl, csAuthInfo.getToken()); +// InputStream content = response.getEntity().getContent(); + InputStream content = HttpUtils.sendGetReturnResponse(presignUrl, csAuthInfo.getToken()); + //存入minio + return fileSaveToMinIO(content, savePath, fileName); + } catch (Exception e) { + log.error("singleFileDownload error", e); + throw new Exception("singleFileDownload error:" + e); + } } + /** + * 所有文件打包上传 + * + * @param inputs + * @param output + * @param zipFilePath + * @throws Exception + */ + @Override + public void allFileUpload(List inputs, String output, String zipFilePath) { + // 异步提交任务 + CompletableFuture future = CompletableFuture.supplyAsync(() -> { + try { + minioUtil.createZipFromDirectory(bucketName, inputs, output, zipFilePath); + return "执行成功"; + } catch (Exception e) { + throw new RuntimeException("异步提交任务失败", e); + } + }); + // 处理异步结果 + future.thenAcceptAsync(result -> { + log.info("所有文件打包上传成功,result: {}", result); + }).exceptionally(ex -> { + log.error("异步处理任务失败", ex); + return null; + }); + } + + /** * 根据 taskParam(表单定义)和 taskInfo(表单提交值)提取文件字段,构建 hpcBindingFiles 列表 + * * @return hpcBindingFiles 列表,每个元素包含 paramName 和 objectID */ - public static List> buildHpcBindingFiles(JSONObject taskParam, JSONObject taskInfo) { + private List> buildHpcBindingFiles(String taskName, JSONObject taskParam, JSONObject taskInfo) { List> bindingFiles = new ArrayList<>(); // 遍历 schema 的所有字段 for (Map.Entry entry : taskParam.entrySet()) { @@ -196,6 +310,12 @@ public class CSCollectServiceImpl implements CSCollectService { Map binding = new HashMap<>(); binding.put("filename", filename); binding.put("packageID", packageId); + if ("specfem3d_globe".equals(taskName)) { + binding.put("paramName", fieldName); + } else { + binding.put("paramName", "inputFile"); + } + binding.put("packageID", packageId); bindingFiles.add(binding); } } @@ -206,15 +326,16 @@ public class CSCollectServiceImpl implements CSCollectService { /** * 将任务提交参数转换为后端所需格式 - * @param clusterId 集群ID - * @param name 任务名称 + * + * @param clusterId 集群ID + * @param name 任务名称 * @param hpcBindingFiles 文件绑定列表,每个元素包含 paramName 和 objectID - * @param userId 用户ID(原入参中的 objectId 视为用户ID) - * @param ntasks 并行任务数 - * @param nodes 节点数 + * @param userId 用户ID(原入参中的 objectId 视为用户ID) + * @param ntasks 并行任务数 + * @param nodes 节点数 * @return 符合JSON串3结构的JSONObject */ - public JSONObject convertToJobRequest(String taskName,String clusterId, String name, + private JSONObject convertToJobRequest(String taskName, String clusterId, String name, List> hpcBindingFiles, Integer userId, String ntasks, String nodes) { // 根对象 @@ -245,13 +366,12 @@ public class CSCollectServiceImpl implements CSCollectService { if (hpcBindingFiles != null) { for (Map fileInfo : hpcBindingFiles) { JSONObject binding = new JSONObject(); - binding.set("paramName", "inputFile"); - JSONObject resource = new JSONObject(); resource.set("type", "path"); - resource.set("path", "/"+fileInfo.get("filename")); + resource.set("path", "/" + fileInfo.get("filename")); resource.set("packageID", fileInfo.get("packageID")); + binding.set("paramName", fileInfo.get("paramName")); binding.set("resource", resource); bindingArray.set(binding); } @@ -292,21 +412,67 @@ public class CSCollectServiceImpl implements CSCollectService { * @throws Exception */ private Integer queryPackageId(String appCode, Integer userId, String token) throws Exception { - //1.获取packageId - Integer packageId = redisService.getCacheObject(PACKAGE_KEY + appCode); - if (packageId == null) { - //2.重新获取packageId - String url = csUrl + "/jsm/v2/package/create"; + try { + //1.获取packageId + Integer packageId = redisService.getCacheObject(PACKAGE_KEY + appCode); + if (packageId == null) { + //1.重新获取packageId + String url = csUrl + "/jsm/v2/package/create"; + JSONObject params = new JSONObject(); + params.set("userID", userId); + params.set("name", appCode + "test"); + params.set("dataType", "HPCSlurm"); + String res = HttpUtils.sendBodyPostWithToken(url, params, token); + JSONObject resData = returnResDataWithOK(res); + packageId = (Integer) getByPath(resData, "newPackage.packageID"); + + //2.获取需要用到的调度平台 + JSONArray hpcAppCluster = getHpcAppCluster(appCode); + binding(userId, packageId, hpcAppCluster); + redisService.setCacheObject(PACKAGE_KEY + appCode, packageId); + } + return packageId; + } catch (Exception e) { + log.error("queryPackageId error", e); + throw new Exception("queryPackageId error:" + e); + } + } + + /** + * 文件夹关联调度平台 + * + * @param userId + * @param packageID + * @param hpcAppCluster + * @return + */ + private void binding(Integer userId, Integer packageID, JSONArray hpcAppCluster) throws Exception { + try { + String url = csUrl + "/jsm/jobSet/binding"; + + JSONObject info = new JSONObject(); + info.set("operateType", "schedule"); + info.set("packageID", packageID); + info.set("type", "HPCSlurm"); + + List idList = hpcAppCluster.stream() + .map(item -> { + JSONObject obj = (JSONObject) item; + return obj.getStr("id"); + }) + .collect(Collectors.toList()); + info.set("clusterIDs", idList); + JSONObject params = new JSONObject(); params.set("userID", userId); - params.set("name", appCode + "test"); - params.set("dataType", "HPCSlurm"); - String res = HttpUtils.sendBodyPostWithToken(url, params, token); - JSONObject resData = returnResDataWithOK(res); - packageId = (Integer) getByPath(resData, "newPackage.packageID"); - redisService.setCacheObject(PACKAGE_KEY + appCode, packageId); + params.set("info", info); + + String result = HttpUtils.sendPost(url, null, JSONUtil.toJsonStr(params)); + returnResDataWithOK(result); + } catch (Exception e) { + log.error(e.getMessage(), e); + throw e; } - return packageId; } /** @@ -316,14 +482,19 @@ public class CSCollectServiceImpl implements CSCollectService { * @throws Exception */ private CSAuthInfo login() throws Exception { - String url = csUrl + "/jcc-admin/admin/login"; - Dict params = Dict.create().set("username", username).set("password", password); - String result = HttpUtils.sendPost(url, null, JSONUtil.toJsonStr(params)); - JSONObject resData = returnResData(result); - Integer userId = (Integer) getByPath(resData, "jsmUserInfo.data.userID"); - Integer tokenTimeout = (Integer) getByPath(resData, "tokenTimeout"); - String token = (String) getByPath(resData, "token"); - return new CSAuthInfo(userId, token, tokenTimeout); + try { + String url = csUrl + "/jcc-admin/admin/login"; + Dict params = Dict.create().set("username", username).set("password", password); + String result = HttpUtils.sendPost(url, null, JSONUtil.toJsonStr(params)); + JSONObject resData = returnResData(result); + Integer userId = (Integer) getByPath(resData, "jsmUserInfo.data.userID"); + Integer tokenTimeout = (Integer) getByPath(resData, "tokenTimeout"); + String token = (String) getByPath(resData, "token"); + return new CSAuthInfo(userId, token, tokenTimeout); + } catch (Exception e) { + log.error("login error", e); + throw new Exception("login error:" + e); + } } /** @@ -334,15 +505,20 @@ public class CSCollectServiceImpl implements CSCollectService { * @return url */ private String presign(Integer userId, String token, Integer packageId, String fileName) throws Exception { - String url = csUrl + "/jsm/v2/storage/presign"; - String presignParam = "{\"userID\":137,\"info\":{\"type\":\"upload\",\"params\":{\"packageID\":13863,\"path\":\"2.jpg\",\"copyTo\":[],\"copyToPath\":[]}}}"; - JSONObject params = JSONUtil.parseObj(presignParam); - params.set("userID", userId); - params.putByPath("info.params.packageID", packageId); - params.putByPath("info.params.path", fileName); - String res = HttpUtils.sendBodyPostWithToken(url, params, token); - JSONObject resData = returnResDataWithOK(res); - return (String) getByPath(resData, "presignUrl"); + try { + String url = csUrl + "/jsm/v2/storage/presign"; + String presignParam = "{\"userID\":137,\"info\":{\"type\":\"upload\",\"params\":{\"packageID\":13863,\"path\":\"2.jpg\",\"copyTo\":[],\"copyToPath\":[]}}}"; + JSONObject params = JSONUtil.parseObj(presignParam); + params.set("userID", userId); + params.putByPath("info.params.packageID", packageId); + params.putByPath("info.params.path", fileName); + String res = HttpUtils.sendBodyPostWithToken(url, params, token); + JSONObject resData = returnResDataWithOK(res); + return (String) getByPath(resData, "presignUrl"); + }catch (Exception e){ + log.error("presign error", e); + throw new Exception("presign error:" + e); + } } /** @@ -364,97 +540,22 @@ public class CSCollectServiceImpl implements CSCollectService { } - /** - * 任务执行结果查询 - * - * @param taskId - * @throws Exception - */ - @Override - public JSONArray queryTaskResult(String taskId) throws Exception { - try { - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - String url = csUrl + "/jsm/v2/jobs/results"; - JSONObject params = new JSONObject(); - params.set("taskID", taskId); - params.set("type", "HPCSlurm"); - params.set("userID", csAuthInfo.getUserId()); - String res = HttpUtils.sendBodyPostWithToken(url, params, csAuthInfo.getToken()); - JSONObject resData = returnResDataWithOK(res); - return JSONUtil.parseArray(getByPath(resData, "pcmJobData.resultFiles[0].objects")); - }catch (Exception e){ - return null; - } - } - - - /** - * 单个文件下载 - * - * @param objectId - * @param savePath - * @param fileName - * @throws Exception - */ - @Override - public JSONObject singleFileDownload(Integer objectId, String savePath, String fileName) throws Exception { - CSAuthInfo csAuthInfo = getZSAuthInfoFromRedis(); - String url = csUrl + "/jsm/v2/storage/presign"; - String paramStr = "{\"userID\":137,\"info\":{\"type\":\"download\",\"params\":{\"objectID\":70805}}}"; - JSONObject params = JSONUtil.parseObj(paramStr); - params.set("userID", csAuthInfo.getUserId()); - params.putByPath("info.params.objectID", objectId); - //获取下载地址 - String res = HttpUtils.sendBodyPostWithToken(url, params, csAuthInfo.getToken()); - JSONObject resData = returnResDataWithOK(res); - String presignUrl = resData.get("presignUrl").toString(); - //下载文件 -// CloseableHttpResponse response = HttpUtils.sendGetReturnResponse(presignUrl, csAuthInfo.getToken()); -// InputStream content = response.getEntity().getContent(); - InputStream content = HttpUtils.sendGetReturnResponse(presignUrl, csAuthInfo.getToken()); - //存入minio - return fileSaveToMinIO(content, savePath, fileName); - } - - /** - * 所有文件打包上传 - * - * @param inputs - * @param output - * @param zipFilePath - * @throws Exception - */ - @Override - public void allFileUpload(List inputs, String output, String zipFilePath) throws Exception { - // 异步提交任务 - CompletableFuture future = CompletableFuture.supplyAsync(() -> { - try { - minioUtil.createZipFromDirectory(bucketName, inputs,output, zipFilePath); - return "执行成功"; - } catch (Exception e) { - throw new RuntimeException("异步提交任务失败", e); - } - }); - // 处理异步结果 - future.thenAcceptAsync(result -> { - log.info("所有文件打包上传成功,result: {}", result); - }).exceptionally(ex -> { - log.error("异步处理任务失败", ex); - return null; - }); - } - //将内容同步到minio中 private JSONObject fileSaveToMinIO(InputStream inputStream, String savePath, String fileName) throws Exception { - String path = savePath + "/" + fileName; - ByteArrayInputStream byteArrayOutputStream = readAllBytes(inputStream); - minioUtil.uploadObject(bucketName, path, byteArrayOutputStream); + try { + String path = savePath + "/" + fileName; + ByteArrayInputStream byteArrayOutputStream = readAllBytes(inputStream); + minioUtil.uploadObject(bucketName, path, byteArrayOutputStream); - JSONObject result = new JSONObject(); - result.set("download_url", path); - result.set("filename", fileName); - result.set("name", fileName); - return result; + JSONObject result = new JSONObject(); + result.set("download_url", path); + result.set("filename", fileName); + result.set("name", fileName); + return result; + } catch (Exception e) { + log.error("fileSaveToMinIO error", e); + throw new Exception("fileSaveToMinIO error:" + e); + } } /** diff --git a/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/utils/MinioUtil.java b/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/utils/MinioUtil.java index 99dae844..053a0aa8 100644 --- a/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/utils/MinioUtil.java +++ b/ruoyi-modules/management-platform/src/main/java/com/ruoyi/platform/utils/MinioUtil.java @@ -140,6 +140,7 @@ public class MinioUtil { /** * 设置响应头下载文件 + * * @param bucket * @param objectPath * @return @@ -518,10 +519,9 @@ public class MinioUtil { /** * 创建目录的压缩包 */ - public void createZipFromDirectory(String bucket, List inputs,String output, String zipPath) - throws Exception { + public void createZipFromDirectory(String bucket, List inputs, String output, String zipPath) throws Exception { - Path tempFile = Files.createTempFile("zip-", ".tmp"); + Path tempFile = Files.createTempFile("zip-", ".tmp"); try { // 写入临时文件 try (FileOutputStream fos = new FileOutputStream(tempFile.toFile()); @@ -546,6 +546,9 @@ public class MinioUtil { .build() ); } + } catch (Exception e) { + log.error("createZipFromDirectory error!", e); + throw new Exception(e); } finally { // 清理临时文件 Files.deleteIfExists(tempFile); @@ -557,55 +560,55 @@ public class MinioUtil { * @param folderName 压缩包内的文件夹名 * @param sourcePath MinIO中的路径 */ - private void writeBuffer(String folderName,String sourcePath,String bucket,ZipOutputStream zos) { - try { + private void writeBuffer(String folderName, String sourcePath, String bucket, ZipOutputStream zos) { + try { - // 1. 规范化文件夹名 - if (!folderName.endsWith("/")) { - folderName = folderName + "/"; - } - - // 2. 获取目录下的所有文件 - Iterable> items = minioClient.listObjects( - ListObjectsArgs.builder() - .bucket(bucket) - .prefix(sourcePath) - .recursive(true) - .build() - ); - - // 3. 添加文件到ZIP - for (Result result : items) { - Item item = result.get(); - if (item.isDir()) continue; - - String objectName = item.objectName(); - String relativePath = objectName.substring(sourcePath.length() - 1); - - // 使用自定义的文件夹名 - String zipEntryName = folderName + relativePath; - - // ZIP条目 - zos.putNextEntry(new ZipEntry(zipEntryName)); - - try (InputStream is = minioClient.getObject( - GetObjectArgs.builder() - .bucket(bucket) - .object(objectName) - .build())) { - - // 使用缓冲区读写 - byte[] buffer = new byte[8192]; - int bytesRead; - while ((bytesRead = is.read(buffer)) != -1) { - zos.write(buffer, 0, bytesRead); - } - } - - zos.closeEntry(); - } - }catch (Exception e){ - e.printStackTrace(); + // 1. 规范化文件夹名 + if (!folderName.endsWith("/")) { + folderName = folderName + "/"; } + + // 2. 获取目录下的所有文件 + Iterable> items = minioClient.listObjects( + ListObjectsArgs.builder() + .bucket(bucket) + .prefix(sourcePath) + .recursive(true) + .build() + ); + + // 3. 添加文件到ZIP + for (Result result : items) { + Item item = result.get(); + if (item.isDir()) continue; + + String objectName = item.objectName(); + String relativePath = objectName.substring(sourcePath.length() - 1); + + // 使用自定义的文件夹名 + String zipEntryName = folderName + relativePath; + + // ZIP条目 + zos.putNextEntry(new ZipEntry(zipEntryName)); + + try (InputStream is = minioClient.getObject( + GetObjectArgs.builder() + .bucket(bucket) + .object(objectName) + .build())) { + + // 使用缓冲区读写 + byte[] buffer = new byte[8192]; + int bytesRead; + while ((bytesRead = is.read(buffer)) != -1) { + zos.write(buffer, 0, bytesRead); + } + } + + zos.closeEntry(); + } + } catch (Exception e) { + log.error("createZipFromDirectory error!", e); } + } } \ No newline at end of file