Spring中Runnable类Autowire及多文件异步处理优化问询
文件多线程处理优化方案梳理
问题背景
我有一个服务,通过for循环处理多个文件,文件可分为5类,处理流程主要为读取文件、调用REST接口更新数据库,每个类别对应不同的REST接口,且处理速度较慢。
原有单线程代码结构
父类代码
public class Parent { public static final String SUCCESS_MESSAGE = "Successfully processed your files"; private static final String SPRING_PROFILES_ACTIVE = "spring.profiles.active"; private static final String LINE_SEPERATOR = System.getProperty("line.separator"); @Autowired private ChildService1 childservice1; @Autowired private ChildService2 childservice2; @Autowired private ChildService3 childservice3; @Autowired private ChildService4 childservice4; /** * 处理从FTP服务器下载的文件,遍历目录中所有文件并分类处理 * @return 处理完成提示字符串 */ @Override public String processFiles() { for (File file : files) { try { if (file.getName().contains("svalue1")) childservice1.processFile(file); else if (file.getName().contains("svalue2")) childservice2.processFile(file); else if (file.getName().contains("value3")) childservice3.processFile(file); else if (file.getName().contains("value4")) childservice4.processFile(file); } catch (Exception e) { log.error("处理文件{}时出现异常,原因:{}", file.getName(), e.getLocalizedMessage(), e); } } return SUCCESS_MESSAGE; } }
子类代码
@Service @Slf4j public class ChildService { @Autowired private ReadCSVService<ChildServiceCSV> readCSVService; @Autowired private Mapper mapperUtil; @Autowired private ChildRepository repository; /** * 读取文件内容,将每行映射为ChildServiceCSV对象,再转换为REST请求模型,最后调用仓库方法更新数据库 * @param file 待处理文件 */ public void processFile(File file) { Iterator<ChildServiceCSV> iterator = readCSVService.mapToCSVIterator(file.getAbsolutePath(), ChildServiceCSV.class); while (iterator.hasNext()) { ChildServiceCSV csvModel = iterator.next(); RESTModel restModel = mapperUtil.mapCsvModelToMDMModel(csvModel); repository.update(restModel); } } }
优化思路与问题
计划通过ExecutorService提升处理速度,在父类中自动装配线程池,将文件处理任务放到独立线程中执行。但遇到设计问题:子服务的processFile方法需要接收File参数,若让子服务实现Runnable,如何调用带参数的方法?
尝试过给子服务添加带File参数的构造函数,将文件赋值给成员变量后在run方法中调用,但这样无法在父类中自动装配子服务,必须每次创建新实例。期望实现的逻辑如下:
for (File file : files) { if (file.getName().contains("svalue1")) executerService.execute(childservice1.processFile(file)); }
已实现的非自动装配方案
线程池配置
@Bean public ExecutorService getExecuted(){ return Executors.newFixedThreadPool(10); }
父类修改
public class Parent { @Autowired ExecutorService executorService; public String processFiles() { for (File file : files) { if (file.getName().contains("svalue1")) executorService.execute( new ChildService1(file)); //其他子服务同理 } executorService.shutdown(); return SUCCESS_MESSAGE; } }
子类修改
@Service @Slf4j public class ChildService implements Runnable { @Autowired private ReadCSVService<ChildServiceCSV> readCSVService; @Autowired private Mapper mapperUtil; @Autowired private ChildRepository repository; private File file; ChildService (File file){ this.file=file; } /** * 读取文件内容,将每行映射为ChildServiceCSV对象,再转换为REST请求模型,最后调用仓库方法更新数据库 * @param file 待处理文件 */ public void processFile(File file) { Iterator<ChildServiceCSV> iterator = readCSVService.mapToCSVIterator(file.getAbsolutePath(), ChildServiceCSV.class); while (iterator.hasNext()) { ChildServiceCSV csvModel = iterator.next(); RESTModel restModel = mapperUtil.mapCsvModelToMDMModel(csvModel); repository.update(restModel); } } @Override public void run() { this.processFile(file); } }
更新:问题已解决
使用Spring的@EnableAsync和@Async注解实现,既保留自动装配,又实现多线程处理:
线程池配置
@SpringBootApplication @EnableAsync public class CreateTaskConsoleApplication implements CommandLineRunner { @Bean public Executor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(10); executor.setMaxPoolSize(10); executor.setQueueCapacity(500); executor.setThreadNamePrefix("Thread-"); executor.initialize(); return executor; } //.....其他代码 }
子类修改
在processFile方法上添加@Async注解,返回类型改为CompletableFuture<String>:
@Async public CompletableFuture<String> processFile(File file) { //原有处理逻辑 Iterator<ChildServiceCSV> iterator = readCSVService.mapToCSVIterator(file.getAbsolutePath(), ChildServiceCSV.class); while (iterator.hasNext()) { ChildServiceCSV csvModel = iterator.next(); RESTModel restModel = mapperUtil.mapCsvModelToMDMModel(csvModel); repository.update(restModel); } return CompletableFuture.completedFuture("文件处理完成"); }
父类修改
收集所有异步任务的CompletableFuture,等待全部完成后返回:
public String processFiles() { List<CompletableFuture<String>> futuresList = new ArrayList<>(); for (File file : files) { try { if (file.getName().contains("svalue1")) futuresList.add(childservice1.processFile(file)); else if (file.getName().contains("svalue2")) futuresList.add(childservice2.processFile(file)); else if (file.getName().contains("value3")) futuresList.add(childservice3.processFile(file)); else if (file.getName().contains("value4")) futuresList.add(childservice4.processFile(file)); } catch (Exception e) { log.error("处理文件{}时出现异常,原因:{}", file.getName(), e.getLocalizedMessage(), e); } } CompletableFuture.allOf(futuresList.toArray(new CompletableFuture[futuresList.size()])).join(); return SUCCESS_MESSAGE; }
内容的提问来源于stack exchange,提问作者Saransh
相关产品推荐
相关产品推荐

