标注@Async的processEvents函数执行两次仅一次生效问题排查
问题描述
我有一个watchDirectoryService用于监听指定目录的变化,当事件触发时,会基于传入异步processEvents函数的Processor接口实现类运行处理流程。但出现异常:虽两个进程看似在后台运行,但仅最后执行的函数会被触发生效;若注释其中一个异步processEvents函数,未注释的可正常工作。
相关组件:
WatchDirService:负责监听已注册目录集合,事件触发时调用对应实现类执行任务DecryptEncryptedFileProcessor:实现Processor接口,执行解密任务ProcessDecryptedFileProcessor:实现Processor接口,执行文件处理任务
相关代码
WatchDirService 代码
public class WatchDirService { private WatchService watcher; private final Map<WatchKey, Path> keys = new HashMap<>(); private boolean recursive; @SuppressWarnings("unchecked") static <T> WatchEvent<T> cast(WatchEvent<?> event) { return (WatchEvent<T>) event; } /** * Register the given directory with the WatchService * * @param dir * @throws java.io.IOException */ private void register(Path dir) throws IOException { WatchKey key = dir.register(watcher, ENTRY_CREATE, ENTRY_DELETE, ENTRY_MODIFY); log.info("PUTTING KEY :: {}", dir); keys.putIfAbsent(key, dir); log.info("KEY COUNT :: {}", keys.size()); } /** * Register the given directory, and all its sub-directories, with the * WatchService. */ private void registerAll(final Path start) throws IOException { // register directory and sub-directories Files.walkFileTree(start, new SimpleFileVisitor<Path>() { @Override public FileVisitResult preVisitDirectory(Path dir, BasicFileAttributes attrs) throws IOException { register(dir); return FileVisitResult.CONTINUE; } }); } /** * Creates a WatchService and registers the given directory * * @param dir * @param recursive * @throws java.io.IOException */ public void register(Path dir, boolean recursive) throws IOException { this.watcher = FileSystems.getDefault().newWatchService(); this.recursive = recursive; if (recursive) { log.info("Scanning {} ...\n", dir); registerAll(dir); log.info("Done."); } else { log.info("Scanning no recursion {} ...\n", dir); register(dir); } } /** * Process all events for keys queued to the watcher * * @param processor * @throws java.lang.Exception */ @Async public void processEvents(Processor processor) throws Exception { log.info("KEY COUNT AND TYPE :: {} {}", keys.size(), processor.getType()); if (keys.isEmpty()) { usage(); } while (true) { log.info("==PROCESSING=={}",processor.getType()); // wait for key to be signalled WatchKey key; try { key = watcher.take(); } catch (InterruptedException x) { return; } Path dir = keys.get(key); if (dir == null) { log.error("WatchKey not recognized!!"); continue; } log.info("PROCESSING DIR :: {}", dir); for (WatchEvent<?> event : key.pollEvents()) { WatchEvent.Kind kind = event.kind(); // TBD - provide example of how OVERFLOW event is handled if (kind == OVERFLOW) { continue; } // Context for directory entry event is the file name of entry WatchEvent<Path> ev = cast(event); Path name = ev.context(); Path child = dir.resolve(name); // print out event log.info("%{}: %{}\n", event.kind().name(), child); // here you would start batch process processor.run(child.toString()); // if directory is created, and watching recursively, then // register it and its sub-directories if (recursive && (kind == ENTRY_CREATE)) { try { if (Files.isDirectory(child, NOFOLLOW_LINKS)) { registerAll(child); } } catch (IOException x) { // ignore to keep sample readbale } } } // reset key and remove from set if directory no longer accessible boolean valid = key.reset(); if (!valid) { keys.remove(key); // all directories are inaccessible if (keys.isEmpty()) { break; } } } log.info("===ENDING PROCESS FOR {}===", processor.getType()); } private void usage() { log.error("usage: java WatchDir [-r] dir"); System.exit(-1); } }
DecryptEncryptedFileProcessor 代码
@Slf4j @RequiredArgsConstructor public class DecryptEncryptedFileProcessor implements Processor<DecryptFile> { @Getter private final DecryptFile command; @Override public void run(String... params) { if (params[0] == null) { throw new RuntimeException("No parameter found to decrypt file, expecting at least one parameter that contains encrypted file"); } FileDecoder f = new FileDecoder(); f.autoBatchDecrypt(params[0], command.getDecryptedFileLocation(), command.getDecryptedFileExtension(), command.getPrefix()); } // shortened }
ProcessDecryptedFileProcessor 代码
@Slf4j @RequiredArgsConstructor public class ProcessDecryptedFileProcessor implements Processor { // private JobRepository jobRepository; private final SpringBatchConfig config; private ItemProcessor<TxnItem, TransactionRecord> processor; @Override public void run(String... params) { if (params[0] == null) { throw new RuntimeException("No decrypted file passed for processing, expecting at least one parameter which contains a decrypted file"); } Path toPath = new File(params[0]).toPath(); String data = readFromInputStream(toPath); log.info("DECRYPTED DATA :: {}", data); } // shortened }
调用代码
@Autowired public Resources( CypherDirectoryConfiguration cypherConfig, WatchDirService watchService) throws IOException, Exception { this.config = config; this.cypherConfig = cypherConfig; this.watchService = watchService; Path encryptedDir = new File(this.cypherConfig.getEncryptedDir()).toPath(); Path decryptedDir = new File(this.cypherConfig.getDecryptedDir()).toPath(); DecryptFile command = new DecryptFile(this.cypherConfig.getDecryptedDir()); DecryptEncryptedFileProcessor dfp = new DecryptEncryptedFileProcessor(command); ProcessDecryptedFileProcessor pdfp = new ProcessDecryptedFileProcessor(this.config); log.info("====STARTING WATCH ENCRYPTED DIR===="); this.watchService.register(encryptedDir, false); this.watchService.processEvents(dfp); log.info("====STARTING WATCH DECRYPTED DIR===="); this.watchService.register(decryptedDir, false); this.watchService.processEvents(pdfp); }
问题原因
核心问题出在WatchDirService的状态共享和WatchService覆盖上:
- Spring默认将
WatchDirService作为单例Bean,每次调用register(Path dir, boolean recursive)时,都会执行this.watcher = FileSystems.getDefault().newWatchService();,直接覆盖之前的WatchService实例。 - 第一次注册加密目录后,启动的异步
processEvents线程持有旧的WatchService引用;但第二次注册解密目录时,全局watcher变量被替换,导致第一个线程的WatchService失效,只有第二个线程能拿到有效实例。 - 全局
keysMap会被两次注册操作叠加,但只有第二个线程的WatchService能接收事件,因此仅最后一个Processor生效。
解决方案
需要让每个监听任务拥有独立的状态或合理绑定目录与处理器,以下是两种可行方案:
方案1:将WatchDirService改为原型Bean
让Spring为每个监听任务创建独立的WatchDirService实例,彻底避免状态共享:
- 在
WatchDirService类上添加@Scope("prototype")注解,确保每次获取时都是新实例。 - 修改调用代码,为每个目录创建独立的监听实例:
@Autowired private ApplicationContext context; public Resources(CypherDirectoryConfiguration cypherConfig) throws IOException, Exception { this.cypherConfig = cypherConfig; Path encryptedDir = new File(this.cypherConfig.getEncryptedDir()).toPath(); Path decryptedDir = new File(this.cypherConfig.getDecryptedDir()).toPath(); DecryptFile command = new DecryptFile(this.cypherConfig.getDecryptedDir()); DecryptEncryptedFileProcessor dfp = new DecryptEncryptedFileProcessor(command); ProcessDecryptedFileProcessor pdfp = new ProcessDecryptedFileProcessor(this.config); // 为加密目录创建独立的WatchDirService WatchDirService encryptedWatchService = context.getBean(WatchDirService.class); log.info("====STARTING WATCH ENCRYPTED DIR===="); encryptedWatchService.register(encryptedDir, false); encryptedWatchService.processEvents(dfp); // 为解密目录创建独立的WatchDirService WatchDirService decryptedWatchService = context.getBean(WatchDirService.class); log.info("====STARTING WATCH DECRYPTED DIR===="); decryptedWatchService.register(decryptedDir, false); decryptedWatchService.processEvents(pdfp); }
方案2:重构WatchDirService支持多目录绑定
如果希望保持WatchDirService为单例,可重构内部逻辑,绑定目录与对应的Processor:
- 在构造函数中初始化一次WatchService,避免重复创建覆盖。
- 维护
Map<Path, Processor>存储目录与处理器的对应关系。 - 修改
processEvents方法,根据事件触发的目录匹配对应的Processor执行。
重构后的核心代码示例
public class WatchDirService { private final WatchService watcher; private final Map<WatchKey, Path> keys = new HashMap<>(); private final Map<Path, Processor> dirProcessorMap = new HashMap<>(); private boolean recursive; // 构造函数中初始化WatchService public WatchDirService() throws IOException { this.watcher = FileSystems.getDefault().newWatchService(); } // 修改register方法,绑定目录和Processor public void register(Path dir, boolean recursive, Processor processor) throws IOException { this.recursive = recursive; if (recursive) { registerAll(dir); } else { register(dir); } dirProcessorMap.put(dir, processor); } @Async public void processEvents() throws Exception { if (keys.isEmpty()) { usage(); } while (true) { WatchKey key = watcher.take(); Path dir = keys.get(key); if (dir == null) { log.error("WatchKey not recognized!!"); continue; } Processor processor = dirProcessorMap.get(dir); if (processor == null) { log.warn("No processor found for dir: {}", dir); continue; } for (WatchEvent<?> event : key.pollEvents()) { WatchEvent.Kind kind = event.kind(); if (kind == OVERFLOW) { continue; } WatchEvent<Path> ev = cast(event); Path name = ev.context(); Path child = dir.resolve(name); log.info("%{}: %{}\n", event.kind().name(), child); processor.run(child.toString()); if (recursive && (kind == ENTRY_CREATE)) { try { if (Files.isDirectory(child, NOFOLLOW_LINKS)) { registerAll(child); // 子目录继承父目录的Processor dirProcessorMap.put(child, processor); } } catch (IOException x) { // ignore } } } boolean valid = key.reset(); if (!valid) { keys.remove(key); dirProcessorMap.remove(dir); if (keys.isEmpty()) { break; } } } } // 其他方法保持不变... }
修改后的调用代码
@Autowired public Resources(CypherDirectoryConfiguration cypherConfig, WatchDirService watchService) throws IOException, Exception { this.cypherConfig = cypherConfig; this.watchService = watchService; Path encryptedDir = new File(this.cypherConfig.getEncryptedDir()).toPath(); Path decryptedDir = new File(this.cypherConfig.getDecryptedDir()).toPath(); DecryptFile command = new DecryptFile(this.cypherConfig.getDecryptedDir()); DecryptEncryptedFileProcessor dfp = new DecryptEncryptedFileProcessor(command); ProcessDecryptedFileProcessor pdfp = new ProcessDecryptedFileProcessor(this.config); log.info("====STARTING WATCH ENCRYPTED DIR===="); this.watchService.register(encryptedDir, false, dfp); log.info("====STARTING WATCH DECRYPTED DIR===="); this.watchService.register(decryptedDir, false, pdfp); // 启动一次事件处理线程即可 this.watchService.processEvents(); }
内容的提问来源于stack exchange,提问作者DaviesTobi alex
相关产品推荐
相关产品推荐

