我刚刚编写了一些与您的代码非常接近的测试用例:
@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();
}
}
};