diff --git a/applications/source/s3-source/README.adoc b/applications/source/s3-source/README.adoc index 3fcf4e73..9f6ac3af 100644 --- a/applications/source/s3-source/README.adoc +++ b/applications/source/s3-source/README.adoc @@ -75,7 +75,6 @@ $$metadata.store.zookeeper.root$$:: $$Root node - store entries are children of $$s3.common.endpoint-url$$:: $$Optional endpoint url to connect to s3 compatible storage.$$ *($$String$$, default: `$$$$`)* $$s3.common.path-style-access$$:: $$Use path style access.$$ *($$Boolean$$, default: `$$false$$`)* $$s3.supplier.auto-create-local-dir$$:: $$Create or not the local directory.$$ *($$Boolean$$, default: `$$true$$`)* -$$s3.supplier.delay-when-empty$$:: $$Duration of delay when no new files are detected.$$ *($$Duration$$, default: `$$1s$$`)* $$s3.supplier.delete-remote-files$$:: $$Delete or not remote files after processing.$$ *($$Boolean$$, default: `$$false$$`)* $$s3.supplier.filename-pattern$$:: $$The pattern to filter remote files.$$ *($$String$$, default: `$$$$`)* $$s3.supplier.filename-regex$$:: $$The regexp to filter remote files.$$ *($$Pattern$$, default: `$$$$`)* @@ -115,6 +114,6 @@ And for AWS `Stack`: == Examples ``` -java -jar s3-source.jar --s3.remoteDir=/tmp/foo --file.consumer.mode=lines --trigger.fixed-delay=60 +java -jar s3-source.jar --s3.remoteDir=/tmp/foo --file.consumer.mode=lines ``` //end::ref-doc[] diff --git a/functions/common/aws-s3-common/pom.xml b/functions/common/aws-s3-common/pom.xml index c739dd08..5488cefa 100644 --- a/functions/common/aws-s3-common/pom.xml +++ b/functions/common/aws-s3-common/pom.xml @@ -14,7 +14,7 @@ - 2.3.0.RELEASE + 2.3.3.RELEASE 2.2.2.RELEASE diff --git a/functions/common/aws-s3-common/src/main/java/org/springframework/cloud/fn/common/aws/s3/CompatibleStorageAmazonS3Configuration.java b/functions/common/aws-s3-common/src/main/java/org/springframework/cloud/fn/common/aws/s3/CompatibleStorageAmazonS3Configuration.java index efaac3ac..73b66836 100644 --- a/functions/common/aws-s3-common/src/main/java/org/springframework/cloud/fn/common/aws/s3/CompatibleStorageAmazonS3Configuration.java +++ b/functions/common/aws-s3-common/src/main/java/org/springframework/cloud/fn/common/aws/s3/CompatibleStorageAmazonS3Configuration.java @@ -16,30 +16,41 @@ package org.springframework.cloud.fn.common.aws.s3; +import java.net.URI; +import java.net.URISyntaxException; + import com.amazonaws.auth.AWSCredentialsProvider; import com.amazonaws.client.builder.AwsClientBuilder.EndpointConfiguration; import com.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.AmazonS3ClientBuilder; import org.springframework.boot.autoconfigure.AutoConfigureBefore; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.cloud.aws.core.env.ResourceIdResolver; import org.springframework.cloud.aws.core.region.RegionProvider; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.integration.aws.support.S3SessionFactory; +import org.springframework.lang.Nullable; +import org.springframework.util.StringUtils; /** * @author Timo Salm + * @author David Turanski */ @Configuration @EnableConfigurationProperties(AmazonS3Properties.class) @AutoConfigureBefore(AmazonS3Configuration.class) + public class CompatibleStorageAmazonS3Configuration { @Bean @ConditionalOnProperty("s3.common.endpoint-url") - public AmazonS3 compatibleStorageAmazonS3(AWSCredentialsProvider awsCredentialsProvider, RegionProvider regionProvider, - AmazonS3Properties amazonS3Properties) { + public AmazonS3 compatibleStorageAmazonS3(AWSCredentialsProvider awsCredentialsProvider, + RegionProvider regionProvider, + AmazonS3Properties amazonS3Properties) { final AmazonS3ClientBuilder builder = AmazonS3ClientBuilder.standard(); final EndpointConfiguration endpointConfiguration = new EndpointConfiguration( amazonS3Properties.getEndpointUrl(), regionProvider.getRegion().getName()); @@ -49,4 +60,24 @@ public class CompatibleStorageAmazonS3Configuration { .withPathStyleAccessEnabled(amazonS3Properties.isPathStyleAccess()) .build(); } + + @Bean + @ConditionalOnMissingBean + public S3SessionFactory s3SessionFactory(AmazonS3 amazonS3, @Nullable ResourceIdResolver resourceIdResolver, + AmazonS3Properties amazonS3Properties) { + S3SessionFactory s3SessionFactory = new S3SessionFactory(amazonS3, resourceIdResolver); + if (StringUtils.hasText(amazonS3Properties.getEndpointUrl())) { + URI uri; + try { + uri = new URI(amazonS3Properties.getEndpointUrl()); + } + catch (URISyntaxException e) { + throw new IllegalArgumentException(amazonS3Properties.getEndpointUrl() + " is not a valid URI"); + } + + s3SessionFactory.setEndpoint(String.join(":", uri.getHost(), String.valueOf(uri.getPort()))); + } + return s3SessionFactory; + } + } diff --git a/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java b/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java index cb7f50a4..fc499871 100644 --- a/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java +++ b/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java @@ -17,21 +17,22 @@ package org.springframework.cloud.fn.supplier.s3; import java.io.File; +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.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.model.ListObjectsRequest; import com.amazonaws.services.s3.model.S3ObjectSummary; -import reactor.core.publisher.Flux; -import reactor.core.publisher.MonoProcessor; import org.reactivestreams.Publisher; import org.reactivestreams.Subscription; +import reactor.core.publisher.Flux; +import reactor.core.publisher.MonoProcessor; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.aws.core.env.ResourceIdResolver; import org.springframework.cloud.fn.common.file.FileConsumerProperties; import org.springframework.cloud.fn.common.file.FileUtils; import org.springframework.context.annotation.Bean; @@ -42,7 +43,6 @@ import org.springframework.integration.aws.support.S3SessionFactory; import org.springframework.integration.aws.support.filters.S3PersistentAcceptOnceFileListFilter; import org.springframework.integration.aws.support.filters.S3RegexPatternFileListFilter; import org.springframework.integration.aws.support.filters.S3SimplePatternFileListFilter; -import org.springframework.integration.core.GenericSelector; import org.springframework.integration.core.MessageSource; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.endpoint.ReactiveMessageSourceProducer; @@ -69,18 +69,18 @@ public abstract class AwsS3SupplierConfiguration { protected final AmazonS3 amazonS3; - protected final ResourceIdResolver resourceIdResolver; + protected final S3SessionFactory s3SessionFactory; protected final ConcurrentMetadataStore metadataStore; public AwsS3SupplierConfiguration(AwsS3SupplierProperties awsS3SupplierProperties, FileConsumerProperties fileConsumerProperties, AmazonS3 amazonS3, - ResourceIdResolver resourceIdResolver, ConcurrentMetadataStore metadataStore) { + S3SessionFactory s3SessionFactory, ConcurrentMetadataStore metadataStore) { this.awsS3SupplierProperties = awsS3SupplierProperties; this.fileConsumerProperties = fileConsumerProperties; this.amazonS3 = amazonS3; - this.resourceIdResolver = resourceIdResolver; + this.s3SessionFactory = s3SessionFactory; this.metadataStore = metadataStore; } @@ -114,9 +114,9 @@ public abstract class AwsS3SupplierConfiguration { SynchronizingConfiguration(AwsS3SupplierProperties awsS3SupplierProperties, FileConsumerProperties fileConsumerProperties, AmazonS3 amazonS3, - ResourceIdResolver resourceIdResolver, + S3SessionFactory s3SessionFactory, ConcurrentMetadataStore concurrentMetadataStore) { - super(awsS3SupplierProperties, fileConsumerProperties, amazonS3, resourceIdResolver, + super(awsS3SupplierProperties, fileConsumerProperties, amazonS3, s3SessionFactory, concurrentMetadataStore); } @@ -132,7 +132,7 @@ public abstract class AwsS3SupplierConfiguration { @Bean public S3InboundFileSynchronizer s3InboundFileSynchronizer(ChainFileListFilter filter) { - S3SessionFactory s3SessionFactory = new S3SessionFactory(this.amazonS3, this.resourceIdResolver); + S3InboundFileSynchronizer synchronizer = new S3InboundFileSynchronizer(s3SessionFactory); synchronizer.setDeleteRemoteFiles(this.awsS3SupplierProperties.isDeleteRemoteFiles()); synchronizer.setPreserveTimestamp(this.awsS3SupplierProperties.isPreserveTimestamp()); @@ -162,28 +162,26 @@ public abstract class AwsS3SupplierConfiguration { ListOnlyConfiguration(AwsS3SupplierProperties awsS3SupplierProperties, FileConsumerProperties fileConsumerProperties, - AmazonS3 amazonS3, - ResourceIdResolver resourceIdResolver, ConcurrentMetadataStore metadataStore) { - super(awsS3SupplierProperties, fileConsumerProperties, amazonS3, resourceIdResolver, metadataStore); + AmazonS3 amazonS3, S3SessionFactory s3SessionFactory, ConcurrentMetadataStore metadataStore) { + super(awsS3SupplierProperties, fileConsumerProperties, amazonS3, s3SessionFactory, + metadataStore); + } + + private final MonoProcessor downstreamSubscription = MonoProcessor.create(); + + @Bean + public Supplier>> s3Supplier(Publisher> s3SupplierFlow) { + return () -> Flux.from(s3SupplierFlow) + .doOnSubscribe(downstreamSubscription::onNext); } @Bean - public Supplier>> s3Supplier(Publisher> s3SupplierFlow) { - return () -> Flux.from(s3SupplierFlow); + public Publisher> s3SupplierFlow(ReactiveMessageSourceProducer s3ListingProducer) { + return IntegrationFlows.from(s3ListingProducer).split().toReactivePublisher(); } @Bean - public Publisher> s3SupplierFlow(ReactiveMessageSourceProducer s3ListingProducer, - GenericSelector listOnlyFilter) { - return IntegrationFlows - .from(s3ListingProducer) - .split() - .filter(listOnlyFilter) - .toReactivePublisher(); - } - - @Bean - GenericSelector listOnlyFilter() { + Predicate listOnlyFilter() { Predicate predicate = s -> true; if (StringUtils.hasText(this.awsS3SupplierProperties.getFilenamePattern())) { Pattern pattern = Pattern.compile(this.awsS3SupplierProperties.getFilenamePattern()); @@ -204,21 +202,26 @@ public abstract class AwsS3SupplierConfiguration { return result; }); - GenericSelector selector = predicate::test; - - return selector; + return predicate; } @Bean ReactiveMessageSourceProducer s3ListingMessageProducer(AmazonS3 amazonS3, - AwsS3SupplierProperties awsS3SupplierProperties) { + AwsS3SupplierProperties awsS3SupplierProperties, Predicate filter) { ListObjectsRequest listObjectsRequest = new ListObjectsRequest(); listObjectsRequest.setBucketName(awsS3SupplierProperties.getRemoteDir()); return new ReactiveMessageSourceProducer( - (MessageSource>) () -> new GenericMessage<>( - amazonS3.listObjects(listObjectsRequest).getObjectSummaries())); + (MessageSource>) () -> { + List summaryList = amazonS3.listObjects(listObjectsRequest) + .getObjectSummaries().stream() + .filter(filter).collect(Collectors.toList()); + return summaryList.isEmpty() ? null : new GenericMessage<>(summaryList); + }) { + @Override + protected void subscribeToPublisher(Publisher> publisher) { + super.subscribeToPublisher(Flux.from(publisher).delaySubscription(downstreamSubscription)); + } + }; } - } - }