diff --git a/supplier/spring-s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java b/supplier/spring-s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java index 2fa1d1f9..963827a9 100644 --- a/supplier/spring-s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java +++ b/supplier/spring-s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java @@ -20,7 +20,6 @@ import java.util.List; import java.util.function.Predicate; import java.util.function.Supplier; import java.util.regex.Pattern; -import java.util.stream.Collectors; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; @@ -96,7 +95,7 @@ public class AwsS3SupplierConfiguration { } @Bean - ChainFileListFilter filter(ConcurrentMetadataStore metadataStore) { + ChainFileListFilter s3SupplierFileListFilter(ConcurrentMetadataStore metadataStore) { ChainFileListFilter chainFilter = new ChainFileListFilter<>(); if (StringUtils.hasText(this.awsS3SupplierProperties.getFilenamePattern())) { chainFilter @@ -129,8 +128,7 @@ public class AwsS3SupplierConfiguration { } @Bean - S3InboundFileSynchronizer s3InboundFileSynchronizer(ChainFileListFilter filter) { - + S3InboundFileSynchronizer s3InboundFileSynchronizer(ChainFileListFilter s3SupplierFileListFilter) { S3InboundFileSynchronizer synchronizer = new S3InboundFileSynchronizer(this.s3SessionFactory); synchronizer.setDeleteRemoteFiles(this.awsS3SupplierProperties.isDeleteRemoteFiles()); synchronizer.setPreserveTimestamp(this.awsS3SupplierProperties.isPreserveTimestamp()); @@ -138,8 +136,7 @@ public class AwsS3SupplierConfiguration { synchronizer.setRemoteDirectory(remoteDir); synchronizer.setRemoteFileSeparator(this.awsS3SupplierProperties.getRemoteFileSeparator()); synchronizer.setTemporaryFileSuffix(this.awsS3SupplierProperties.getTmpFileSuffix()); - synchronizer.setFilter(filter); - + synchronizer.setFilter(s3SupplierFileListFilter); return synchronizer; } @@ -178,8 +175,8 @@ public class AwsS3SupplierConfiguration { } @Bean - Publisher> s3SupplierFlow(ReactiveMessageSourceProducer s3ListingProducer) { - return IntegrationFlow.from(s3ListingProducer).split().toReactivePublisher(true); + Publisher> s3SupplierFlow(ReactiveMessageSourceProducer s3ListingMessageProducer) { + return IntegrationFlow.from(s3ListingMessageProducer).split().toReactivePublisher(true); } @Bean @@ -210,13 +207,13 @@ public class AwsS3SupplierConfiguration { @Bean ReactiveMessageSourceProducer s3ListingMessageProducer(S3Client amazonS3, ObjectMapper objectMapper, - AwsS3SupplierProperties awsS3SupplierProperties, Predicate filter) { + AwsS3SupplierProperties awsS3SupplierProperties, Predicate listOnlyFilter) { return new ReactiveMessageSourceProducer((MessageSource>) () -> { List summaryList = amazonS3 .listObjects(ListObjectsRequest.builder().bucket(awsS3SupplierProperties.getRemoteDir()).build()) .contents() .stream() - .filter(filter) + .filter(listOnlyFilter) .map((s3Object) -> { try { return objectMapper.writeValueAsString(s3Object.toBuilder()); @@ -225,7 +222,7 @@ public class AwsS3SupplierConfiguration { throw new RuntimeException(ex); } }) - .collect(Collectors.toList()); + .toList(); return summaryList.isEmpty() ? null : new GenericMessage<>(summaryList); }); }