【问题标题】:WatchServiceDirectoryScanner not pulling new files created after first pollWatchServiceDirectoryScanner 不提取第一次轮询后创建的新文件
【发布时间】:2016-09-11 19:02:38
【问题描述】:

我正在从 Spring Integration 读取 FileReadingMessageSource 中的根目录以检索正在进行的文件创建。场景是根目录下可能有多个子目录持续存在。 SI 4.3.1 中的 WatchServiceDirectoryScanner 用于拾取在任何新子目录中创建的任何文件。

@Bean
public MessageSource<File> fileReadingMessageSource() {

    CompositeFileListFilter<File> filters = new CompositeFileListFilter<>();
    filters.addFilter(new SimplePatternFileListFilter("pattern*"));
    //filters.addFilter(new LastModifiedFileListFilter());

    FileReadingMessageSource fileSource = new FileReadingMessageSource();

    String filePath = "root-directory";

    fileSource.setDirectory(new File(filePath));
    fileSource.setFilter(filters);
    fileSource.setUseWatchService(true);
    fileSource.setWatchEvents(FileReadingMessageSource.WatchEventType.CREATE,FileReadingMessageSource.WatchEventType.MODIFY,FileReadingMessageSource.WatchEventType.DELETE);

    return fileSource;
}

@Bean
public IntegrationFlow readDirectoryFlow() {

    return IntegrationFlows.from(
            fileReadingMessageSource(), 
            e -> e.poller(Pollers.cron("*/5 * * * * *")))
            .channel(fileInputChannel())
            .handle(tailerRestart)
            .handle(System.out::println)
            .get();
}

在第一次轮询时,所有匹配模式的文件都可通过消息资源获得,但如果稍后在任何新子目录中创建任何新文件,则消息资源无法选择新的模式匹配文件。

我在日志中看到以下 DEBUG 消息

DEBUG SourcePollingChannelAdapter - 轮询期间未收到任何消息,返回“false”

可能缺少什么?

【问题讨论】:

    标签: spring spring-integration


    【解决方案1】:

    我刚刚编写了一些与您的代码非常接近的测试用例:

        @Bean
        public MessageSource<File> fileReadingMessageSource() {
            CompositeFileListFilter<File> filters = new CompositeFileListFilter<>();
            filters.addFilter(new SimplePatternFileListFilter("*.watch"));
    
            FileReadingMessageSource fileSource = new FileReadingMessageSource();
            fileSource.setDirectory(tmpDir.getRoot());
            fileSource.setFilter(filters);
            fileSource.setUseWatchService(true);
            fileSource.setWatchEvents(FileReadingMessageSource.WatchEventType.CREATE,
                    FileReadingMessageSource.WatchEventType.MODIFY,
                    FileReadingMessageSource.WatchEventType.DELETE);
            return fileSource;
        }
    
        @Bean
        public IntegrationFlow readDirectoryFlow() {
            return IntegrationFlows
                    .from(fileReadingMessageSource(),
                            e -> e.poller(p -> p.cron("*/1 * * * * *")))
                    .handle(System.out::println)
                    .get();
        }
    

    测试代码如下:

    @ClassRule
    public static final TemporaryFolder tmpDir = new TemporaryFolder();
    
    @Test
    public void testWatchServiceMessageSource() throws Exception {
        File newFolder1 = tmpDir.newFolder();
        FileOutputStream file = new FileOutputStream(new File(newFolder1, "foo.watch"));
        file.write(("foo").getBytes());
        file.flush();
        file.close();
    
        File newFolder2 = tmpDir.newFolder();
        file = new FileOutputStream(new File(newFolder2, "bar.watch"));
        file.write(("bar").getBytes());
        file.flush();
        file.close();
    
        file = new FileOutputStream(new File(tmpDir.getRoot(), "root.watch"));
        file.write(("root").getBytes());
        file.flush();
        file.close();
    
        Thread.sleep(10000);
    }
    

    我有这些日志:

    GenericMessage [payload=C:\Users\abilan\AppData\Local\Temp\junit7602962373770028652\junit7776799219532481336\foo.watch, headers={id=50d44197-e0af-708a-6b61-2a2cfeec68da, timestamp=1473686655061}]
    GenericMessage [payload=C:\Users\abilan\AppData\Local\Temp\junit7602962373770028652\junit813088196038861528\bar.watch, headers={id=8d80c853-19b6-f667-7950-d6de49d509ab, timestamp=1473686656062}]
    GenericMessage [payload=C:\Users\abilan\AppData\Local\Temp\junit7602962373770028652\root.watch, headers={id=e585203b-41dc-cadb-6a36-4c9009a34701, timestamp=1473686657063}]
    

    每秒。

    不知道你的问题出在哪里...

    您不需要.channel(fileInputChannel())。它将在 ednpoints 之间自动创建。

    配置:

    .handle(tailerRestart)
    .handle(System.out::println)
    

    你应该确定tailerRestart 会返回一些东西。虽然,根据我们的其他讨论,它没有:

    @ServiceActivator
    public void restartTailer(File input) throws Exception {
        tailFileProducer.stop();
        tailFileProducer.setFile(input);
        tailFileProducer.start();
    }
    

    更新

    经过一些私人调查,我们发现问题在于 Spring Cloud Stream 基础架构多次调用 FileReadingMessageSource.start(),导致重新实例化内部 WatchService 对象。

    FileReadingMessageSource.start() 必须固定为幂等:https://jira.spring.io/browse/INT-4108

    Spring Cloud Stream 已在版本1.1:https://github.com/spring-cloud/spring-cloud-stream/issues/525 中修复。

    解决方法就像确保FileReadingMessageSource.start() 只被调用一次:

    FileReadingMessageSource fileSource = new FileReadingMessageSource() {
    
       private final AtomicBoolean running = new AtomicBoolean();
    
       @Override
       public void start() {
          if (!this.running.getAndSet(true)) {
             super.start();
          }
       }
    
       @Override
       public void stop() {
          if (this.running.getAndSet(false)) {
             super.stop();
          }
       }
    
    };
    

    【讨论】:

    • 感谢您测试代码。我修改了 restartTailer 服务激活器以返回文件。我觉得 WatchServiceDirectoryScanner 在 FileReadingMessageSource 中只填充一次 Queue toBeReceived 的问题。如果稍后添加文件/子目录,则无法在新子目录中看到这些文件的消息。我们在 WatchServiceDirectoryScanner 中是否有无限轮询来监视来自 Java 7 WatchService 的所有事件?
    • 请在我的回答中找到UPDATE
    猜你喜欢
    • 2020-03-05
    • 1970-01-01
    • 2015-12-30
    • 1970-01-01
    • 2021-08-28
    • 1970-01-01
    • 2019-10-24
    • 2013-11-09
    • 1970-01-01
    相关资源
    最近更新 更多