Fix File Supplier for lifecycle

SO: https://stackoverflow.com/questions/69802199/how-to-use-watchservicedirectoryscanner-with-spring-cloud-stream-file-supplier

The `FileReadingMessageSource` has to be started if `useWatchService()` is used.
This commit is contained in:
Artem Bilan
2021-12-09 17:12:36 -05:00
committed by Soby Chacko
parent 9828120f34
commit 7e54bd2d8a
2 changed files with 12 additions and 1 deletions

View File

@@ -85,7 +85,8 @@ public class FileSupplierConfiguration {
monoSink.success(this.fileMessageSource.receive())))
.subscribeOn(Schedulers.boundedElastic())
.repeatWhenEmpty(it -> it.delayElements(this.fileSupplierProperties.getDelayWhenEmpty()))
.repeat();
.repeat()
.doOnSubscribe(s -> this.fileMessageSource.start());
}
@Bean

View File

@@ -24,7 +24,10 @@ import org.junit.jupiter.api.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.integration.file.FileHeaders;
import org.springframework.integration.file.FileReadingMessageSource;
import org.springframework.messaging.Message;
import static org.assertj.core.api.Assertions.assertThat;
@@ -36,6 +39,10 @@ import static org.assertj.core.api.Assertions.assertThat;
*/
public class DefaultFileSupplierTests extends AbstractFileSupplierTests {
@Autowired
@Qualifier("fileMessageSource")
private FileReadingMessageSource fileMessageSource;
@Test
public void testBasicFlow() throws IOException {
@@ -74,5 +81,8 @@ public class DefaultFileSupplierTests extends AbstractFileSupplierTests {
.verifyLater();
Files.write(tempFile, "testing".getBytes());
stepVerifier.verify();
assertThat(this.fileMessageSource.isRunning()).isTrue();
}
}