You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

标注@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覆盖上:

  1. Spring默认将WatchDirService作为单例Bean,每次调用register(Path dir, boolean recursive)时,都会执行this.watcher = FileSystems.getDefault().newWatchService();,直接覆盖之前的WatchService实例。
  2. 第一次注册加密目录后,启动的异步processEvents线程持有旧的WatchService引用;但第二次注册解密目录时,全局watcher变量被替换,导致第一个线程的WatchService失效,只有第二个线程能拿到有效实例。
  3. 全局keysMap会被两次注册操作叠加,但只有第二个线程的WatchService能接收事件,因此仅最后一个Processor生效。

解决方案

需要让每个监听任务拥有独立的状态或合理绑定目录与处理器,以下是两种可行方案:

方案1:将WatchDirService改为原型Bean

让Spring为每个监听任务创建独立的WatchDirService实例,彻底避免状态共享:

  1. 在WatchDirService类上添加@Scope("prototype")注解,确保每次获取时都是新实例。
  2. 修改调用代码,为每个目录创建独立的监听实例:
@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:

  1. 在构造函数中初始化一次WatchService,避免重复创建覆盖。
  2. 维护Map<Path, Processor>存储目录与处理器的对应关系。
  3. 修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 04:35:54