From fd3253dc385fb8458b2e9fc9692c60941b655e2e Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 25 Sep 2020 13:27:13 -0400 Subject: [PATCH] Defer subscription to the S3 until downstream one It turns out that endpoint for S3 Source is starting to consume the `Flux` earlier than an actual subscription happens for the whole flow downstream * Introduce a `MonoProcessor downstreamSubscription` to fulfill from the downstream subscription and handle from the `delaySubscription()` on the `Flux` for source data --- .../s3/AwsS3SupplierConfiguration.java | 21 ++++++++++++++----- 1 file changed, 16 insertions(+), 5 deletions(-) diff --git a/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java b/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java index ed01b4c0..cb7f50a4 100644 --- a/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java +++ b/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java @@ -24,8 +24,10 @@ import java.util.regex.Pattern; import com.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.model.ListObjectsRequest; import com.amazonaws.services.s3.model.S3ObjectSummary; -import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; +import reactor.core.publisher.MonoProcessor; +import org.reactivestreams.Publisher; +import org.reactivestreams.Subscription; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; @@ -86,9 +88,11 @@ public abstract class AwsS3SupplierConfiguration { @ConditionalOnProperty(prefix = "s3.supplier", name = "list-only", havingValue = "false", matchIfMissing = true) static class SynchronizingConfiguration extends AwsS3SupplierConfiguration { + private final MonoProcessor downstreamSubscription = MonoProcessor.create(); + @Bean - public Supplier>> s3Supplier(Publisher> s3SupplierFlow) { - return () -> Flux.from(s3SupplierFlow); + public Supplier>> s3Supplier(Publisher> s3SupplierFlow) { + return () -> Flux.from(s3SupplierFlow).doOnSubscribe(this.downstreamSubscription::onNext); } @Bean @@ -118,8 +122,11 @@ public abstract class AwsS3SupplierConfiguration { @Bean public Publisher> s3SupplierFlow(MessageSource s3MessageSource) { - return FileUtils.enhanceFlowForReadingMode(IntegrationFlows - .from(IntegrationReactiveUtils.messageSourceToFlux(s3MessageSource)), fileConsumerProperties) + return FileUtils.enhanceFlowForReadingMode( + IntegrationFlows.from( + IntegrationReactiveUtils.messageSourceToFlux(s3MessageSource) + .delaySubscription(this.downstreamSubscription)), + fileConsumerProperties) .toReactivePublisher(); } @@ -146,11 +153,13 @@ public abstract class AwsS3SupplierConfiguration { s3MessageSource.setAutoCreateLocalDirectory(this.awsS3SupplierProperties.isAutoCreateLocalDir()); return s3MessageSource; } + } @Configuration @ConditionalOnProperty(prefix = "s3.supplier", name = "list-only", havingValue = "true") static class ListOnlyConfiguration extends AwsS3SupplierConfiguration { + ListOnlyConfiguration(AwsS3SupplierProperties awsS3SupplierProperties, FileConsumerProperties fileConsumerProperties, AmazonS3 amazonS3, @@ -209,5 +218,7 @@ public abstract class AwsS3SupplierConfiguration { (MessageSource>) () -> new GenericMessage<>( amazonS3.listObjects(listObjectsRequest).getObjectSummaries())); } + } + }