From 946b5969dad9704d4efaaf95b770adce597b47eb Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 7 Mar 2024 14:51:36 -0500 Subject: [PATCH] Various refactoring and improvements * Remove `MessageChannel` beans whenever their really don't need to be exposed into the target application context. Use local property definition instead in the configuration class * Make Twitter `Consumer` configurations conditional on their required properties to avoid auto-configuration for those bean which are not going to be used in the target application --- consumer/spring-twitter-consumer/README.adoc | 14 +++--- ...itterFriendshipsConsumerConfiguration.java | 5 ++- .../TwitterFriendshipsConsumerProperties.java | 13 ------ .../TwitterMessageConsumerConfiguration.java | 4 +- .../TwitterMessageConsumerProperties.java | 2 - .../TwitterUpdateConsumerConfiguration.java | 2 +- ...tterUpdateConsumerConfigurationTests.java} | 5 ++- .../AggregatorFunctionConfiguration.java | 24 ++++------ .../TwitterTrendFunctionConfiguration.java | 25 ++++++----- .../trend/TwitterTrendFunctionTests.java | 7 ++- .../syslog/SyslogSupplierConfiguration.java | 44 +++++++------------ .../TwitterStreamSupplierConfiguration.java | 30 ++++--------- 12 files changed, 69 insertions(+), 106 deletions(-) rename consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/{TwitterUpdateSinkFunctionConfigurationTests.java => TwitterUpdateConsumerConfigurationTests.java} (96%) diff --git a/consumer/spring-twitter-consumer/README.adoc b/consumer/spring-twitter-consumer/README.adoc index 06ea7386..4b60bcdc 100644 --- a/consumer/spring-twitter-consumer/README.adoc +++ b/consumer/spring-twitter-consumer/README.adoc @@ -20,15 +20,15 @@ You can find details for the Update API here: https://developer.twitter.com/en/d The `TwitterUpdateConsumerConfiguration` auto-configuration provides this beans: -- `Consumer updateStatus` - if you have an `StatusUpdate` instance you can use the `updateStatus` to apply it. +- `Consumer twitterUpdateStatusConsumer` - if you have an `StatusUpdate` instance you can use the `updateStatus` to apply it. - `Function, StatusUpdate> messageToStatusUpdateFunction` - function that converts a `Message` text into a `StatusUpdate` instance using the `TwitterUpdateConsumerProperties` properties. -- `Consumer> twitterStatusUpdateConsumer` - composes `messageToStatusUpdateFunction` and `updateStatus` to update the twitter status from Message text. +- `Consumer> twitterUpdateStatusConsumer` - composes `messageToStatusUpdateFunction` and `updateStatus` to update the twitter status from Message text. Note: the Message content is expected to be in text format. Consider using the `byteArrayTextToString` utility `Function`. -You can use `twitterStatusUpdateConsumer` as a qualifier when injecting. +You can use `twitterUpdateStatusConsumer` as a qualifier when injecting. ### 1.2 Configuration Options @@ -58,11 +58,11 @@ SpEL expressions are used to compute the request parameters from the input messa The `TwitterMessageConsumerConfiguration` auto-configuration provides this beans: -- `Consumer> sendDirectMessageConsumer` +- `Consumer> twitterSendMessageConsumer` Note: the Message content is expected to be in text format. Consider using the `byteArrayTextToString` utility `Function`. -You can use `twitterStatusUpdateConsumer` as a qualifier when injecting. +You can use `twitterSendMessageConsumer` as a qualifier when injecting. ### 2.2 Configuration Options @@ -94,11 +94,11 @@ Every operation type has its own parameters. The `TwitterFriendshipsConsumerConfiguration` auto-configuration provides this beans: -- `Consumer> friendshipConsumer` +- `Consumer> twitterFriendshipConsumer` Note: the Message content is expected to be in text format. Consider using the `byteArrayTextToString` utility `Function`. -You can use `friendshipConsumer` as a qualifier when injecting. +You can use `twitterFriendshipConsumer` as a qualifier when injecting. ### 3.2 Configuration Options diff --git a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerConfiguration.java b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerConfiguration.java index f65157a5..0128b9ef 100644 --- a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerConfiguration.java +++ b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerConfiguration.java @@ -22,6 +22,7 @@ import twitter4j.Twitter; import twitter4j.TwitterException; import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.fn.common.twitter.TwitterConnectionConfiguration; import org.springframework.context.annotation.Bean; @@ -35,11 +36,13 @@ import org.springframework.messaging.Message; */ @AutoConfiguration(after = TwitterConnectionConfiguration.class) @EnableConfigurationProperties(TwitterFriendshipsConsumerProperties.class) +@ConditionalOnExpression("environment['twitter.friendships.update.user-id'] !='' or environment['twitter.friendships.update.screen-name'] != ''") public class TwitterFriendshipsConsumerConfiguration { @Bean @SuppressWarnings("Duplicates") - public Consumer> friendshipConsumer(TwitterFriendshipsConsumerProperties properties, Twitter twitter) { + public Consumer> twitterFriendshipConsumer(TwitterFriendshipsConsumerProperties properties, + Twitter twitter) { return (message) -> { try { TwitterFriendshipsConsumerProperties.OperationType type = properties.getType() diff --git a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerProperties.java b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerProperties.java index 97079a6f..6126b840 100644 --- a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerProperties.java +++ b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/friendship/TwitterFriendshipsConsumerProperties.java @@ -16,15 +16,11 @@ package org.springframework.cloud.fn.consumer.twitter.friendship; -import jakarta.validation.constraints.AssertTrue; -import jakarta.validation.constraints.NotNull; - import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.expression.Expression; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.expression.ValueExpression; import org.springframework.stereotype.Component; -import org.springframework.validation.annotation.Validated; /** * The Twitter friendships properties. @@ -34,7 +30,6 @@ import org.springframework.validation.annotation.Validated; */ @Component @ConfigurationProperties("twitter.friendships.update") -@Validated public class TwitterFriendshipsConsumerProperties { public enum OperationType { @@ -101,17 +96,11 @@ public class TwitterFriendshipsConsumerProperties { return this.update; } - @AssertTrue(message = "Either userId or screenName must be provided") - public boolean isUserProvided() { - return this.userId != null || this.screenName != null; - } - public static class Create { /** * The ID of the user to follow (boolean). */ - @NotNull private Expression follow = new ValueExpression<>(true); public Expression getFollow() { @@ -129,13 +118,11 @@ public class TwitterFriendshipsConsumerProperties { /** * Enable/disable device notifications from the target user. */ - @NotNull private Expression device = new ValueExpression<>(true); /** * Enable/disable Retweets from the target user. */ - @NotNull private Expression retweets = new ValueExpression<>(true); public Expression getDevice() { diff --git a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerConfiguration.java b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerConfiguration.java index eb8be7b4..20e0cc8b 100644 --- a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerConfiguration.java +++ b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerConfiguration.java @@ -24,6 +24,7 @@ import twitter4j.Twitter; import twitter4j.TwitterException; import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.fn.common.twitter.TwitterConnectionConfiguration; import org.springframework.context.annotation.Bean; @@ -37,12 +38,13 @@ import org.springframework.messaging.Message; */ @AutoConfiguration(after = TwitterConnectionConfiguration.class) @EnableConfigurationProperties(TwitterMessageConsumerProperties.class) +@ConditionalOnExpression("environment['twitter.message.update.user-id'] !='' or environment['twitter.message.update.screen-name'] != ''") public class TwitterMessageConsumerConfiguration { private static final Log LOGGER = LogFactory.getLog(TwitterMessageConsumerConfiguration.class); @Bean - public Consumer> sendDirectMessageConsumer(TwitterMessageConsumerProperties messageProperties, + public Consumer> twitterSendMessageConsumer(TwitterMessageConsumerProperties messageProperties, Twitter twitter) { return (message) -> { diff --git a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerProperties.java b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerProperties.java index c161dba6..c2b68bda 100644 --- a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerProperties.java +++ b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/message/TwitterMessageConsumerProperties.java @@ -19,7 +19,6 @@ package org.springframework.cloud.fn.consumer.twitter.message; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.expression.Expression; import org.springframework.expression.spel.standard.SpelExpressionParser; -import org.springframework.validation.annotation.Validated; /** * The Twitter messages properties. @@ -27,7 +26,6 @@ import org.springframework.validation.annotation.Validated; * @author Christian Tzolov */ @ConfigurationProperties("twitter.message.update") -@Validated public class TwitterMessageConsumerProperties { private static final Expression DEFAULT_EXPRESSION = new SpelExpressionParser().parseExpression("payload"); diff --git a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfiguration.java b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfiguration.java index 01a0dbdb..bafebe61 100644 --- a/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfiguration.java +++ b/consumer/spring-twitter-consumer/src/main/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfiguration.java @@ -45,7 +45,7 @@ public class TwitterUpdateConsumerConfiguration { private static final Log LOGGER = LogFactory.getLog(TwitterUpdateConsumerConfiguration.class); @Bean - public Consumer updateStatus(Twitter twitter) { + public Consumer twitterUpdateStatusConsumer(Twitter twitter) { return (statusUpdate) -> { try { Status status = twitter.updateStatus(statusUpdate); diff --git a/consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java b/consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfigurationTests.java similarity index 96% rename from consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java rename to consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfigurationTests.java index bd0279f6..99056c92 100644 --- a/consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java +++ b/consumer/spring-twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateConsumerConfigurationTests.java @@ -39,7 +39,7 @@ import static org.mockito.Mockito.verify; * @author Christian Tzolov * @author Artem Bilan */ -public class TwitterUpdateSinkFunctionConfigurationTests { +public class TwitterUpdateConsumerConfigurationTests { private static final ExpressionParser EXPRESSION_PARSER = new SpelExpressionParser(); @@ -47,7 +47,8 @@ public class TwitterUpdateSinkFunctionConfigurationTests { public void testStatusUpdateConsumer() throws TwitterException { Twitter twitter = mock(Twitter.class); - Consumer statusUpdateConsumer = new TwitterUpdateConsumerConfiguration().updateStatus(twitter); + Consumer statusUpdateConsumer = new TwitterUpdateConsumerConfiguration() + .twitterUpdateStatusConsumer(twitter); StatusUpdate statusUpdateQuery = new StatusUpdate("Hello World"); statusUpdateConsumer.accept(statusUpdateQuery); diff --git a/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java b/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java index 70becbfc..032b4596 100644 --- a/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java +++ b/function/spring-aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java @@ -23,7 +23,6 @@ import reactor.core.publisher.Flux; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -46,7 +45,6 @@ import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.support.MessageBuilder; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageChannel; /** * The auto-configuration for aggregator function. @@ -58,6 +56,8 @@ import org.springframework.messaging.MessageChannel; @EnableConfigurationProperties(AggregatorFunctionProperties.class) public class AggregatorFunctionConfiguration { + private final FluxMessageChannel outputChannel = new FluxMessageChannel(); + @Autowired private AggregatorFunctionProperties properties; @@ -65,28 +65,22 @@ public class AggregatorFunctionConfiguration { private BeanFactory beanFactory; @Bean - public Function>, Flux>> aggregatorFunction(FluxMessageChannel inputChannel, - FluxMessageChannel outputChannel) { - return (input) -> Flux.from(outputChannel) - .doOnRequest((request) -> inputChannel.subscribeTo(input.map(( + public Function>, Flux>> aggregatorFunction(FluxMessageChannel aggregatorInputChannel) { + return (input) -> Flux.from(this.outputChannel) + .doOnRequest((request) -> aggregatorInputChannel.subscribeTo(input.map(( inputMessage) -> MessageBuilder.fromMessage(inputMessage).removeHeader("kafka_consumer").build()))); } @Bean - public FluxMessageChannel inputChannel() { + public FluxMessageChannel aggregatorInputChannel() { return new FluxMessageChannel(); } @Bean - public FluxMessageChannel outputChannel() { - return new FluxMessageChannel(); - } - - @Bean - @ServiceActivator(inputChannel = "inputChannel") + @ServiceActivator(inputChannel = "aggregatorInputChannel") public AggregatorFactoryBean aggregator(@Nullable CorrelationStrategy correlationStrategy, @Nullable ReleaseStrategy releaseStrategy, @Nullable MessageGroupProcessor messageGroupProcessor, - @Nullable MessageGroupStore messageStore, @Qualifier("outputChannel") MessageChannel outputChannel, + @Nullable MessageGroupStore messageStore, @Nullable ComponentCustomizer aggregatorCustomizer) { AggregatorFactoryBean aggregator = new AggregatorFactoryBean(); @@ -112,7 +106,7 @@ public class AggregatorFunctionConfiguration { if (messageStore != null) { aggregator.setMessageStore(messageStore); } - aggregator.setOutputChannel(outputChannel); + aggregator.setOutputChannel(this.outputChannel); if (aggregatorCustomizer != null) { aggregatorCustomizer.customize(aggregator); diff --git a/function/spring-twitter-function/src/main/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionConfiguration.java b/function/spring-twitter-function/src/main/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionConfiguration.java index 6f2f5705..00e5160e 100644 --- a/function/spring-twitter-function/src/main/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionConfiguration.java +++ b/function/spring-twitter-function/src/main/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionConfiguration.java @@ -38,6 +38,7 @@ import org.springframework.messaging.Message; * Auto-configuration for Twitter Trend function. * * @author Christian Tzolov + * @author Artem Bilan */ @ConditionalOnProperty(prefix = "twitter.trend", name = "trend-query-type") @AutoConfiguration(after = TwitterConnectionConfiguration.class) @@ -47,7 +48,17 @@ public class TwitterTrendFunctionConfiguration { private static final Log LOGGER = LogFactory.getLog(TwitterTrendFunctionConfiguration.class); @Bean - public Function, Trends> trend(TwitterTrendFunctionProperties properties, Twitter twitter) { + public Function, Message> twitterTrendFunction(TwitterTrendFunctionProperties properties, + Twitter twitter, Function> managedJson) { + + Function, Trends> trendsFunction = trendsFunction(properties, twitter); + Function, List> closestOrAvailableTrends = closestOrAvailableTrends(properties, twitter); + + return (properties.getTrendQueryType() == TwitterTrendFunctionProperties.TrendQueryType.trend) + ? trendsFunction.andThen(managedJson) : closestOrAvailableTrends.andThen(managedJson); + } + + private Function, Trends> trendsFunction(TwitterTrendFunctionProperties properties, Twitter twitter) { return (message) -> { try { int woeid = properties.getLocationId().getValue(message, int.class); @@ -60,8 +71,7 @@ public class TwitterTrendFunctionConfiguration { }; } - @Bean - public Function, List> closestOrAvailableTrends(TwitterTrendFunctionProperties properties, + private Function, List> closestOrAvailableTrends(TwitterTrendFunctionProperties properties, Twitter twitter) { return (message) -> { @@ -82,13 +92,4 @@ public class TwitterTrendFunctionConfiguration { }; } - @Bean - public Function, Message> twitterTrendFunction(Function> managedJson, - Function, Trends> trend, TwitterTrendFunctionProperties properties, - Function, List> closestOrAvailableTrends) { - - return (properties.getTrendQueryType() == TwitterTrendFunctionProperties.TrendQueryType.trend) - ? trend.andThen(managedJson) : closestOrAvailableTrends.andThen(managedJson); - } - } diff --git a/function/spring-twitter-function/src/test/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionTests.java b/function/spring-twitter-function/src/test/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionTests.java index 1fa07062..00088d50 100644 --- a/function/spring-twitter-function/src/test/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionTests.java +++ b/function/spring-twitter-function/src/test/java/org/springframework/cloud/fn/twitter/trend/TwitterTrendFunctionTests.java @@ -49,6 +49,7 @@ import static org.mockserver.verify.VerificationTimes.once; /** * @author Christian Tzolov + * @author Artem Bilan */ @SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = { "twitter.connection.consumerKey=consumerKey666", @@ -57,8 +58,6 @@ import static org.mockserver.verify.VerificationTimes.once; @DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) public abstract class TwitterTrendFunctionTests { - private static final String MOCK_SERVER_IP = "127.0.0.1"; - private static ClientAndServer mockServer; private static MockServerClient mockClient; @@ -71,7 +70,7 @@ public abstract class TwitterTrendFunctionTests { @BeforeAll public static void startServer() { mockServer = ClientAndServer.startClientAndServer(); - mockClient = new MockServerClient(MOCK_SERVER_IP, mockServer.getPort()); + mockClient = new MockServerClient("localhost", mockServer.getPort()); trendsRequest = setExpectation( request().withMethod("GET").withPath("/trends/place.json").withQueryStringParameter("id", "2972")); @@ -116,7 +115,7 @@ public abstract class TwitterTrendFunctionTests { Function mockedConfiguration = toConfigurationBuilder .andThen(new TwitterTestUtils() - .mockTwitterUrls(String.format("http://%s:%s", MOCK_SERVER_IP, mockServer.getPort()))); + .mockTwitterUrls(String.format("http://localhost:" + mockServer.getPort()))); return mockedConfiguration.apply(properties).build(); } diff --git a/supplier/spring-syslog-supplier/src/main/java/org/springframework/cloud/fn/supplier/syslog/SyslogSupplierConfiguration.java b/supplier/spring-syslog-supplier/src/main/java/org/springframework/cloud/fn/supplier/syslog/SyslogSupplierConfiguration.java index 02f5ed26..885b311d 100644 --- a/supplier/spring-syslog-supplier/src/main/java/org/springframework/cloud/fn/supplier/syslog/SyslogSupplierConfiguration.java +++ b/supplier/spring-syslog-supplier/src/main/java/org/springframework/cloud/fn/supplier/syslog/SyslogSupplierConfiguration.java @@ -55,17 +55,14 @@ public class SyslogSupplierConfiguration { @Autowired private SyslogSupplierProperties properties; - @Bean - public FluxMessageChannel syslogInputChannel() { - return new FluxMessageChannel(); - } + private final FluxMessageChannel syslogInputChannel = new FluxMessageChannel(); @Bean - public Supplier>> syslogSupplier(FluxMessageChannel syslogInputChannel, + public Supplier>> syslogSupplier( ObjectProvider udpAdapterProvider, ObjectProvider tcpAdapterProvider) { - return () -> Flux.from(syslogInputChannel).doOnSubscribe((subscription) -> { + return () -> Flux.from(this.syslogInputChannel).doOnSubscribe((subscription) -> { UdpSyslogReceivingChannelAdapter udpAdapter = udpAdapterProvider.getIfAvailable(); TcpSyslogReceivingChannelAdapter tcpAdapter = tcpAdapterProvider.getIfAvailable(); if (udpAdapter != null) { @@ -79,25 +76,19 @@ public class SyslogSupplierConfiguration { @Bean @ConditionalOnProperty(name = "syslog.supplier.protocol", havingValue = "udp") - public UdpSyslogReceivingChannelAdapter udpAdapter(MessageConverter syslogConverter, - FluxMessageChannel syslogInputChannel) { - - return createUdpAdapter(syslogConverter, syslogInputChannel); + public UdpSyslogReceivingChannelAdapter udpAdapter(MessageConverter syslogConverter) { + return createUdpAdapter(syslogConverter); } @Bean @ConditionalOnProperty(name = "syslog.supplier.protocol", havingValue = "both") - public UdpSyslogReceivingChannelAdapter udpBothAdapter(MessageConverter syslogConverter, - FluxMessageChannel syslogInputChannel) { - - return createUdpAdapter(syslogConverter, syslogInputChannel); + public UdpSyslogReceivingChannelAdapter udpBothAdapter(MessageConverter syslogConverter) { + return createUdpAdapter(syslogConverter); } - private UdpSyslogReceivingChannelAdapter createUdpAdapter(MessageConverter syslogConverter, - FluxMessageChannel syslogInputChannel) { - + private UdpSyslogReceivingChannelAdapter createUdpAdapter(MessageConverter syslogConverter) { UdpSyslogReceivingChannelAdapter adapter = new UdpSyslogReceivingChannelAdapter(); - setAdapterProperties(adapter, syslogConverter, syslogInputChannel); + setAdapterProperties(adapter, syslogConverter); return adapter; } @@ -105,18 +96,18 @@ public class SyslogSupplierConfiguration { @ConditionalOnProperty(name = "syslog.supplier.protocol", havingValue = "tcp", matchIfMissing = true) public TcpSyslogReceivingChannelAdapter tcpAdapter( @Qualifier("syslogSupplierConnectionFactory") AbstractServerConnectionFactory connectionFactory, - MessageConverter syslogConverter, FluxMessageChannel syslogInputChannel) { + MessageConverter syslogConverter) { - return createTcpAdapter(connectionFactory, syslogConverter, syslogInputChannel); + return createTcpAdapter(connectionFactory, syslogConverter); } @Bean @ConditionalOnProperty(name = "syslog.supplier.protocol", havingValue = "both") public TcpSyslogReceivingChannelAdapter tcpBothAdapter( @Qualifier("syslogSupplierConnectionFactory") AbstractServerConnectionFactory connectionFactory, - MessageConverter syslogConverter, FluxMessageChannel syslogInputChannel) { + MessageConverter syslogConverter) { - return createTcpAdapter(connectionFactory, syslogConverter, syslogInputChannel); + return createTcpAdapter(connectionFactory, syslogConverter); } @Bean @@ -130,20 +121,19 @@ public class SyslogSupplierConfiguration { } private TcpSyslogReceivingChannelAdapter createTcpAdapter(AbstractServerConnectionFactory connectionFactory, - MessageConverter syslogConverter, FluxMessageChannel syslogInputChannel) { + MessageConverter syslogConverter) { TcpSyslogReceivingChannelAdapter adapter = new TcpSyslogReceivingChannelAdapter(); adapter.setConnectionFactory(connectionFactory); - setAdapterProperties(adapter, syslogConverter, syslogInputChannel); + setAdapterProperties(adapter, syslogConverter); return adapter; } - private void setAdapterProperties(SyslogReceivingChannelAdapterSupport adapter, MessageConverter syslogConverter, - FluxMessageChannel syslogInputChannel) { + private void setAdapterProperties(SyslogReceivingChannelAdapterSupport adapter, MessageConverter syslogConverter) { adapter.setPort(this.properties.getPort()); adapter.setConverter(syslogConverter); - adapter.setOutputChannel(syslogInputChannel); + adapter.setOutputChannel(this.syslogInputChannel); adapter.setAutoStartup(false); } diff --git a/supplier/spring-twitter-supplier/src/main/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierConfiguration.java b/supplier/spring-twitter-supplier/src/main/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierConfiguration.java index 330dd291..a45db5b2 100644 --- a/supplier/spring-twitter-supplier/src/main/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierConfiguration.java +++ b/supplier/spring-twitter-supplier/src/main/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierConfiguration.java @@ -53,14 +53,10 @@ public class TwitterStreamSupplierConfiguration { private static final Log LOGGER = LogFactory.getLog(TwitterStreamSupplierConfiguration.class); - @Bean - public FluxMessageChannel twitterStatusInputChannel() { - return new FluxMessageChannel(); - } + private final FluxMessageChannel twitterStatusInputChannel = new FluxMessageChannel(); @Bean - public StatusListener twitterStatusListener(FluxMessageChannel twitterStatusInputChannel, - TwitterStream twitterStream, ObjectMapper objectMapper) { + public StatusListener twitterStatusListener(TwitterStream twitterStream, ObjectMapper objectMapper) { StatusListener statusListener = new StatusListener() { @@ -94,7 +90,7 @@ public class TwitterStreamSupplierConfiguration { Message message = MessageBuilder.withPayload(json.getBytes()) .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON_VALUE) .build(); - twitterStatusInputChannel.send(message); + TwitterStreamSupplierConfiguration.this.twitterStatusInputChannel.send(message); } catch (JsonProcessingException ex) { String errorMessage = "Status to JSON conversion error!"; @@ -116,23 +112,15 @@ public class TwitterStreamSupplierConfiguration { @Bean public Supplier>> twitterStreamSupplier(TwitterStream twitterStream, - FluxMessageChannel twitterStatusInputChannel, TwitterStreamSupplierProperties streamProperties) { + TwitterStreamSupplierProperties streamProperties) { - return () -> Flux.from(twitterStatusInputChannel).doOnSubscribe((subscription) -> { + return () -> Flux.from(this.twitterStatusInputChannel).doOnSubscribe((subscription) -> { try { switch (streamProperties.getType()) { - case filter -> { - twitterStream.filter(streamProperties.getFilter().toFilterQuery()); - } - case sample -> { - twitterStream.sample(); - } - case firehose -> { - twitterStream.firehose(streamProperties.getFilter().getCount()); - } - case link -> { - twitterStream.links(streamProperties.getFilter().getCount()); - } + case filter -> twitterStream.filter(streamProperties.getFilter().toFilterQuery()); + case sample -> twitterStream.sample(); + case firehose -> twitterStream.firehose(streamProperties.getFilter().getCount()); + case link -> twitterStream.links(streamProperties.getFilter().getCount()); default -> throw new IllegalArgumentException("Unknown stream type:" + streamProperties.getType()); } }