如何在Spring Cloud Stream文件供应商中使用WatchServiceDirectoryScanner?
问题根因
你遇到的异常由自定义配置的逻辑和时机错误直接导致,核心原因有两个:
FileInboundChannelAdapterSpec是构建器对象,其内部FileReadingMessageSource的生命周期本该由对应的入站通道端点统一管理。你在postProcessBeforeInitialization阶段调用spec.get()提前获取了尚未完成构建的MessageSource实例,破坏了原有构建流程,导致WatchServiceDirectoryScanner的启动钩子没有被正确注册到端点的生命周期回调中。- 当你配置
useWatchService(true)后,WatchService的start()方法本应在Spring上下文完成初始化后、端点启动阶段被自动调用。但你提前暴露的MessageSource如果被其他Bean提前触发了接收逻辑,或者端点的生命周期回调缺失,就会出现扫描器未启动就被调用的异常。
解决方案
方案1(推荐):直接在Bean定义时完成配置,放弃BeanPostProcessor自定义的方式
这是最符合Spring Integration规范的用法,避免破坏原有构建和生命周期流程,示例代码如下:
@Bean public FileInboundChannelAdapterSpec fileInboundChannelAdapterSpec(YourProperties properties, FileListFilter<File> fileFilters) { return Files.inboundAdapter(properties.getDirectory()) .autoCreateDirectory(false) .useWatchService(true) .watchEvents(WatchEventType.CREATE, WatchEventType.MODIFY) .preventDuplicates(false) .nioLocker() .filter(fileFilters); }
方案2:如果必须全局自定义所有FileInboundChannelAdapterSpec,调整处理逻辑
不要在BeanPostProcessor中提前调用spec.get()操作内部对象,所有配置仅通过Spec的API设置即可,同时调整处理阶段到postProcessAfterInitialization:
@Bean public BeanPostProcessor inboundFileAdaptorCustomizer() { return new BeanPostProcessor() { @Override public Object postProcessAfterInitialization(Object bean, String beanName) { if (bean instanceof FileInboundChannelAdapterSpec spec) { spec.autoCreateDirectory(false) .useWatchService(true) .watchEvents(WatchEventType.CREATE, WatchEventType.MODIFY) .preventDuplicates(false) .nioLocker() .filter(fileFilters()); } return bean; } }; }
临时兜底方案
如果上述方案都无法适配你的场景,你可以额外注册一个高优先级的SmartLifecycleBean,在上下文启动时主动触发扫描器的启动:
@Bean public SmartLifecycle startWatchService(FileReadingMessageSource fileReadingMessageSource) { return new SmartLifecycle() { private boolean running = false; @Override public int getPhase() { // 优先级高于默认的端点启动阶段 return Integer.MAX_VALUE - 100; } @Override public void start() { fileReadingMessageSource.start(); running = true; } @Override public void stop() { fileReadingMessageSource.stop(); running = false; } @Override public boolean isRunning() { return running; } // 其余默认方法按需求实现即可 }; }
内容的提问来源于stack exchange,提问作者E-Riz
相关产品推荐
相关产品推荐

