多线程实现主位置同步上传+多位置异步写入文件技术咨询
实现多线程异步上传至多个副存储位置
看起来你需要完成主存储同步上传+多副存储异步并行上传的需求,咱们可以通过线程池来高效实现多路径并行写入,同时优化文件读取逻辑避免重复IO开销。下面是具体的实现方案:
核心思路
- 优先同步完成主存储的文件上传,确保主存储写入成功后再触发副存储的异步任务
- 使用线程池管理副存储的上传线程,控制并发数避免资源耗尽
- 针对不同文件大小优化读取逻辑:小文件可缓存至内存复用,大文件则为每个线程独立打开文件流(或用零拷贝技术)
完整代码实现
1. 初始化存储位置列表(延续你的代码片段)
// 初始化主副存储位置列表 List<FileUploadMultiLocator> fileUploadList = new ArrayList<>(); // 主存储位置 FileUploadMultiLocator primaryStorage = new FileUploadMultiLocator("fileserver0", new File("/primary/storage/file.txt")); // 副存储位置1 FileUploadMultiLocator secondaryStorage1 = new FileUploadMultiLocator("fileserver1", new File("/secondary/storage1/file.txt")); // 副存储位置2 FileUploadMultiLocator secondaryStorage2 = new FileUploadMultiLocator("fileserver2", new File("/secondary/storage2/file.txt")); fileUploadList.add(primaryStorage); fileUploadList.add(secondaryStorage1); fileUploadList.add(secondaryStorage2);
2. 配置线程池
根据副存储的数量配置合适的线程池,平衡并发性能和资源占用:
// 线程池核心线程数设为副存储数量,最大线程数不超过CPU核心数的2倍 int secondaryStorageCount = fileUploadList.size() - 1; // 减去主存储 ThreadPoolExecutor uploadExecutor = new ThreadPoolExecutor( secondaryStorageCount, Math.min(secondaryStorageCount * 2, Runtime.getRuntime().availableProcessors() * 2), 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), Executors.defaultThreadFactory(), new ThreadPoolExecutor.AbortPolicy() // 任务满时拒绝新任务,可根据业务调整 );
3. 同步处理主存储上传
先确保主存储上传成功,再启动副存储的异步任务:
String sourceFilePath = "/path/to/your/source/file.txt"; // 同步上传至主存储 try (InputStream primaryInputStream = new FileInputStream(sourceFilePath)) { uploadToStorage(primaryInputStream, fileUploadList.get(0)); System.out.println("主存储上传完成"); } catch (IOException e) { System.err.println("主存储上传失败,终止副存储任务:" + e.getMessage()); uploadExecutor.shutdownNow(); // 终止所有未执行的副存储任务 return; }
4. 多线程异步上传至副存储
为每个副存储提交独立的上传任务,注意流的正确关闭:
// 遍历副存储位置(跳过第一个主存储) for (int i = 1; i < fileUploadList.size(); i++) { FileUploadMultiLocator secondaryLocator = fileUploadList.get(i); // 提交异步任务 uploadExecutor.submit(() -> { try (InputStream secondaryInputStream = new FileInputStream(sourceFilePath)) { uploadToStorage(secondaryInputStream, secondaryLocator); System.out.println("副存储 " + secondaryLocator.getServerName() + " 上传完成"); } catch (IOException e) { System.err.println("副存储 " + secondaryLocator.getServerName() + " 上传失败:" + e.getMessage()); } }); } // 等待所有副存储任务完成(可选,根据业务需求决定是否阻塞等待) uploadExecutor.shutdown(); try { if (!uploadExecutor.awaitTermination(30, TimeUnit.MINUTES)) { System.err.println("部分副存储任务超时,强制终止"); uploadExecutor.shutdownNow(); } } catch (InterruptedException e) { uploadExecutor.shutdownNow(); Thread.currentThread().interrupt(); }
5. 通用上传逻辑封装
把写入存储的逻辑抽成通用方法,方便后续扩展不同存储类型(如FTP、云存储):
/** * 通用文件上传方法,根据存储位置信息写入文件 */ private void uploadToStorage(InputStream inputStream, FileUploadMultiLocator locator) throws IOException { // 这里可以根据locator的类型实现不同的存储写入逻辑 // 示例:写入本地文件系统 try (OutputStream outputStream = new FileOutputStream(locator.getFilePath())) { byte[] buffer = new byte[8192]; // 8KB缓冲区,平衡IO性能和内存占用 int bytesRead; while ((bytesRead = inputStream.read(buffer)) != -1) { outputStream.write(buffer, 0, bytesRead); } outputStream.flush(); // 确保所有数据写入磁盘 } }
关键优化点
- 大文件处理:如果是GB级大文件,建议使用
FileChannel的transferTo/transferFrom实现零拷贝,避免内存溢出:// 大文件零拷贝示例 try (FileChannel sourceChannel = new FileInputStream(sourceFilePath).getChannel(); FileChannel targetChannel = new FileOutputStream(locator.getFilePath()).getChannel()) { sourceChannel.transferTo(0, sourceChannel.size(), targetChannel); } - 异常隔离:每个副存储任务独立捕获异常,避免一个任务失败导致其他任务中断
- 资源清理:必须使用
try-with-resources自动关闭流,线程池使用后必须shutdown,避免资源泄漏
内容的提问来源于stack exchange,提问作者Harisingh Rajput
相关产品推荐
相关产品推荐

