From b695aa51b5eba296c31cf218b268d9b22caf12c6 Mon Sep 17 00:00:00 2001 From: hesamghiasi Date: Tue, 20 Sep 2022 07:16:47 +0430 Subject: [PATCH] SEDA architecture with si and camel (#12588) * adding getCurrentTime method to TimeAgoCalculatorUnitTest in order to always return the same time and avoid problems related to reading date from local system. adding two methods to TimeAgoCalculator to always return the same date as the current date in order to avoid problems related to reading current time from local host. One of these two methods accepts time zone * adding getCurrentTime method to TimeAgoCalculatorUnitTest in order to always return the same time and avoid problems related to reading date from local system. adding two methods to TimeAgoCalculator to always return the same date as the current date in order to avoid problems related to reading current time from local host. One of these two methods accepts time zone correcting some formattings adding comments in code to clarify adding of getCurrentTime methods * reverting changes to ZuulConfig * adding seda package to design-patterns-architectural to demonstrate the implementation of SEDA architectural style in spring integration and apache camel in solving word count problem. * changing version of apache camel to be compatible with java 8 * changing toList() to Collectors.toList() in order to be compatile with java 8 * fixing package names in test classes * fixing string for scanning integration gateway * resolving comments on pull request * resolving comments on pull request * renaming java file of a class * changing variable name * add a line before assertion * some formatting of fluent API * using transform instead of service activator * change name of class * BAEL-5636 Fixed compilation error and applied Baeldung formatter Co-authored-by: bipster --- .../design-patterns-architectural/pom.xml | 24 ++++++ .../seda/apachecamel/WordCountRoute.java | 48 +++++++++++ .../ChannelConfiguration.java | 56 +++++++++++++ .../IntegrationConfiguration.java | 81 +++++++++++++++++++ .../TaskExecutorConfiguration.java | 61 ++++++++++++++ .../seda/springintegration/TestGateway.java | 12 +++ .../ApacheCamelSedaIntegrationTest.java | 48 +++++++++++ .../SpringIntegrationSedaIntegrationTest.java | 50 ++++++++++++ 8 files changed, 380 insertions(+) create mode 100644 patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/apachecamel/WordCountRoute.java create mode 100644 patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/ChannelConfiguration.java create mode 100644 patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/IntegrationConfiguration.java create mode 100644 patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TaskExecutorConfiguration.java create mode 100644 patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TestGateway.java create mode 100644 patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/apachecamel/ApacheCamelSedaIntegrationTest.java create mode 100644 patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/springintegration/SpringIntegrationSedaIntegrationTest.java diff --git a/patterns-modules/design-patterns-architectural/pom.xml b/patterns-modules/design-patterns-architectural/pom.xml index 5cfe90867b..8f579ddfa2 100644 --- a/patterns-modules/design-patterns-architectural/pom.xml +++ b/patterns-modules/design-patterns-architectural/pom.xml @@ -41,6 +41,26 @@ ${mysql-connector.version} jar + + org.springframework.boot + spring-boot-starter-integration + ${spring-boot-starter-integration.version} + + + org.springframework.integration + spring-integration-test + ${spring-integration-test.version} + + + org.apache.camel + camel-core + ${camel-core.version} + + + org.apache.camel + camel-test-junit5 + ${camel-test-junit5.version} + @@ -48,6 +68,10 @@ 6.0.6 2.5.3 3.3.0 + 2.7.2 + 5.5.14 + 3.14.0 + 3.14.0 \ No newline at end of file diff --git a/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/apachecamel/WordCountRoute.java b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/apachecamel/WordCountRoute.java new file mode 100644 index 0000000000..c84a755426 --- /dev/null +++ b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/apachecamel/WordCountRoute.java @@ -0,0 +1,48 @@ +package com.baeldung.seda.apachecamel; + +import java.util.List; +import java.util.function.Function; +import java.util.stream.Collectors; + +import org.apache.camel.Exchange; +import org.apache.camel.builder.ExpressionBuilder; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.processor.aggregate.AbstractListAggregationStrategy; + +public class WordCountRoute extends RouteBuilder { + + public static final String receiveTextUri = "seda:receiveText?concurrentConsumers=5"; + public static final String splitWordsUri = "seda:splitWords?concurrentConsumers=5"; + public static final String toLowerCaseUri = "seda:toLowerCase?concurrentConsumers=5"; + public static final String countWordsUri = "seda:countWords?concurrentConsumers=5"; + public static final String returnResponse = "mock:result"; + + @Override + public void configure() throws Exception { + + from(receiveTextUri).to(splitWordsUri); + + from(splitWordsUri).transform(ExpressionBuilder.bodyExpression(s -> s.toString() + .split(" "))) + .to(toLowerCaseUri); + + from(toLowerCaseUri).split(body(), new StringListAggregationStrategy()) + .transform(ExpressionBuilder.bodyExpression(body -> body.toString() + .toLowerCase())) + .end() + .to(countWordsUri); + + from(countWordsUri).transform(ExpressionBuilder.bodyExpression(List.class, body -> body.stream() + .collect(Collectors.groupingBy(Function.identity(), Collectors.counting())))) + .to(returnResponse); + } +} + +class StringListAggregationStrategy extends AbstractListAggregationStrategy { + @Override + public String getValue(Exchange exchange) { + return exchange.getIn() + .getBody(String.class); + } + +} diff --git a/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/ChannelConfiguration.java b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/ChannelConfiguration.java new file mode 100644 index 0000000000..027ea6a83d --- /dev/null +++ b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/ChannelConfiguration.java @@ -0,0 +1,56 @@ +package com.baeldung.seda.springintegration; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.task.TaskExecutor; +import org.springframework.integration.dsl.MessageChannels; +import org.springframework.messaging.MessageChannel; + +@Configuration +public class ChannelConfiguration { + + private final TaskExecutor receiveTextChannelThreadPool; + private final TaskExecutor splitWordsChannelThreadPool; + private final TaskExecutor toLowerCaseChannelThreadPool; + private final TaskExecutor countWordsChannelThreadPool; + private final TaskExecutor returnResponseChannelThreadPool; + + public ChannelConfiguration(TaskExecutor receiveTextChannelThreadPool, TaskExecutor splitWordsChannelThreadPool, TaskExecutor toLowerCaseChannelThreadPool, TaskExecutor countWordsChannelThreadPool, TaskExecutor returnResponseChannelThreadPool) { + this.receiveTextChannelThreadPool = receiveTextChannelThreadPool; + this.splitWordsChannelThreadPool = splitWordsChannelThreadPool; + this.toLowerCaseChannelThreadPool = toLowerCaseChannelThreadPool; + this.countWordsChannelThreadPool = countWordsChannelThreadPool; + this.returnResponseChannelThreadPool = returnResponseChannelThreadPool; + } + + @Bean(name = "receiveTextChannel") + public MessageChannel getReceiveTextChannel() { + return MessageChannels.executor("receive-text", receiveTextChannelThreadPool) + .get(); + } + + @Bean(name = "splitWordsChannel") + public MessageChannel getSplitWordsChannel() { + return MessageChannels.executor("split-words", splitWordsChannelThreadPool) + .get(); + } + + @Bean(name = "toLowerCaseChannel") + public MessageChannel getToLowerCaseChannel() { + return MessageChannels.executor("to-lower-case", toLowerCaseChannelThreadPool) + .get(); + } + + @Bean(name = "countWordsChannel") + public MessageChannel getCountWordsChannel() { + return MessageChannels.executor("count-words", countWordsChannelThreadPool) + .get(); + } + + @Bean(name = "returnResponseChannel") + public MessageChannel getReturnResponseChannel() { + return MessageChannels.executor("return-response", returnResponseChannelThreadPool) + .get(); + } + +} diff --git a/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/IntegrationConfiguration.java b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/IntegrationConfiguration.java new file mode 100644 index 0000000000..7df4ee51c8 --- /dev/null +++ b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/IntegrationConfiguration.java @@ -0,0 +1,81 @@ +package com.baeldung.seda.springintegration; + +import java.util.List; +import java.util.Map; +import java.util.function.Function; +import java.util.stream.Collectors; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.aggregator.MessageGroupProcessor; +import org.springframework.integration.aggregator.ReleaseStrategy; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.dsl.IntegrationFlows; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.support.MessageBuilder; + +@Configuration +@EnableIntegration +public class IntegrationConfiguration { + + private final MessageChannel receiveTextChannel; + private final MessageChannel splitWordsChannel; + private final MessageChannel toLowerCaseChannel; + private final MessageChannel countWordsChannel; + private final MessageChannel returnResponseChannel; + + private final Function splitWordsFunction = sentence -> sentence.split(" "); + private final Function, Map> convertArrayListToCountMap = list -> list.stream() + .collect(Collectors.groupingBy(Function.identity(), Collectors.counting())); + private final Function toLowerCase = String::toLowerCase; + private final MessageGroupProcessor buildMessageWithListPayload = messageGroup -> MessageBuilder.withPayload(messageGroup.streamMessages() + .map(Message::getPayload) + .collect(Collectors.toList())) + .build(); + private final ReleaseStrategy listSizeReached = r -> r.size() == r.getSequenceSize(); + + public IntegrationConfiguration(MessageChannel receiveTextChannel, MessageChannel splitWordsChannel, MessageChannel toLowerCaseChannel, MessageChannel countWordsChannel, MessageChannel returnResponseChannel) { + this.receiveTextChannel = receiveTextChannel; + this.splitWordsChannel = splitWordsChannel; + this.toLowerCaseChannel = toLowerCaseChannel; + this.countWordsChannel = countWordsChannel; + this.returnResponseChannel = returnResponseChannel; + } + + @Bean + public IntegrationFlow receiveText() { + return IntegrationFlows.from(receiveTextChannel) + .channel(splitWordsChannel) + .get(); + } + + @Bean + public IntegrationFlow splitWords() { + return IntegrationFlows.from(splitWordsChannel) + .transform(splitWordsFunction) + .channel(toLowerCaseChannel) + .get(); + } + + @Bean + public IntegrationFlow toLowerCase() { + return IntegrationFlows.from(toLowerCaseChannel) + .split() + .transform(toLowerCase) + .aggregate(aggregatorSpec -> aggregatorSpec.releaseStrategy(listSizeReached) + .outputProcessor(buildMessageWithListPayload)) + .channel(countWordsChannel) + .get(); + } + + @Bean + public IntegrationFlow countWords() { + return IntegrationFlows.from(countWordsChannel) + .transform(convertArrayListToCountMap) + .channel(returnResponseChannel) + .get(); + } + +} diff --git a/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TaskExecutorConfiguration.java b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TaskExecutorConfiguration.java new file mode 100644 index 0000000000..0596d61eb1 --- /dev/null +++ b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TaskExecutorConfiguration.java @@ -0,0 +1,61 @@ +package com.baeldung.seda.springintegration; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.task.TaskExecutor; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +@Configuration +public class TaskExecutorConfiguration { + + @Bean("receiveTextChannelThreadPool") + public TaskExecutor receiveTextChannelThreadPool() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(1); + executor.setMaxPoolSize(5); + executor.setThreadNamePrefix("receive-text-channel-thread-pool"); + executor.initialize(); + return executor; + } + + @Bean("splitWordsChannelThreadPool") + public TaskExecutor splitWordsChannelThreadPool() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(1); + executor.setMaxPoolSize(5); + executor.setThreadNamePrefix("split-words-channel-thread-pool"); + executor.initialize(); + return executor; + } + + @Bean("toLowerCaseChannelThreadPool") + public TaskExecutor toLowerCaseChannelThreadPool() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(1); + executor.setMaxPoolSize(5); + executor.setThreadNamePrefix("tto-lower-case-channel-thread-pool"); + executor.initialize(); + return executor; + } + + @Bean("countWordsChannelThreadPool") + public TaskExecutor countWordsChannelThreadPool() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(1); + executor.setMaxPoolSize(5); + executor.setThreadNamePrefix("count-words-channel-thread-pool"); + executor.initialize(); + return executor; + } + + @Bean("returnResponseChannelThreadPool") + public TaskExecutor returnResponseChannelThreadPool() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(1); + executor.setMaxPoolSize(5); + executor.setThreadNamePrefix("return-response-channel-thread-pool"); + executor.initialize(); + return executor; + } + +} diff --git a/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TestGateway.java b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TestGateway.java new file mode 100644 index 0000000000..7729477bb7 --- /dev/null +++ b/patterns/design-patterns-architectural/src/main/java/com/baeldung/seda/springintegration/TestGateway.java @@ -0,0 +1,12 @@ +package com.baeldung.seda.springintegration; + +import java.util.Map; + +import org.springframework.integration.annotation.Gateway; +import org.springframework.integration.annotation.MessagingGateway; + +@MessagingGateway +public interface TestGateway { + @Gateway(requestChannel = "receiveTextChannel", replyChannel = "returnResponseChannel") + public Map countWords(String test); +} diff --git a/patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/apachecamel/ApacheCamelSedaIntegrationTest.java b/patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/apachecamel/ApacheCamelSedaIntegrationTest.java new file mode 100644 index 0000000000..778165c356 --- /dev/null +++ b/patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/apachecamel/ApacheCamelSedaIntegrationTest.java @@ -0,0 +1,48 @@ +package com.baeldung.seda.apachecamel; + +import java.util.HashMap; +import java.util.Map; + +import org.apache.camel.RoutesBuilder; +import org.apache.camel.test.junit5.CamelTestSupport; +import org.junit.jupiter.api.Test; + +import com.baeldung.seda.apachecamel.WordCountRoute; + +public class ApacheCamelSedaIntegrationTest extends CamelTestSupport { + + @Test + public void givenTextWithCapitalAndSmallCaseAndWithoutDuplicateWords_whenSendingTextToInputUri_thenWordCountReturnedAsMap() throws InterruptedException { + Map expected = new HashMap<>(); + expected.put("my", 1L); + expected.put("name", 1L); + expected.put("is", 1L); + expected.put("hesam", 1L); + getMockEndpoint(WordCountRoute.returnResponse).expectedBodiesReceived(expected); + template.sendBody(WordCountRoute.receiveTextUri, "My name is Hesam"); + + assertMockEndpointsSatisfied(); + } + + @Test + public void givenTextWithDuplicateWords_whenSendingTextToInputUri_thenWordCountReturnedAsMap() throws InterruptedException { + Map expected = new HashMap<>(); + expected.put("the", 3L); + expected.put("dog", 1L); + expected.put("chased", 1L); + expected.put("rabbit", 1L); + expected.put("into", 1L); + expected.put("jungle", 1L); + getMockEndpoint(WordCountRoute.returnResponse).expectedBodiesReceived(expected); + template.sendBody(WordCountRoute.receiveTextUri, "the dog chased the rabbit into the jungle"); + + assertMockEndpointsSatisfied(); + } + + @Override + protected RoutesBuilder createRouteBuilder() throws Exception { + RoutesBuilder wordCountRoute = new WordCountRoute(); + return wordCountRoute; + } + +} diff --git a/patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/springintegration/SpringIntegrationSedaIntegrationTest.java b/patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/springintegration/SpringIntegrationSedaIntegrationTest.java new file mode 100644 index 0000000000..0ca45aff95 --- /dev/null +++ b/patterns/design-patterns-architectural/src/test/java/com/baeldung/seda/springintegration/SpringIntegrationSedaIntegrationTest.java @@ -0,0 +1,50 @@ +package com.baeldung.seda.springintegration; + +import java.util.HashMap; +import java.util.Map; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.integration.annotation.IntegrationComponentScan; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit.jupiter.SpringExtension; +import org.springframework.test.context.support.AnnotationConfigContextLoader; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(loader = AnnotationConfigContextLoader.class, classes = { TaskExecutorConfiguration.class, ChannelConfiguration.class, IntegrationConfiguration.class }) +@EnableIntegration +@IntegrationComponentScan(basePackages = { "com.baeldung.seda.springintegration" }) +public class SpringIntegrationSedaIntegrationTest { + + @Autowired + TestGateway testGateway; + + @Test + void givenTextWithCapitalAndSmallCaseAndWithoutDuplicateWords_whenCallingCountWordOnGateway_thenWordCountReturnedAsMap() { + Map actual = testGateway.countWords("My name is Hesam"); + Map expected = new HashMap<>(); + expected.put("my", 1L); + expected.put("name", 1L); + expected.put("is", 1L); + expected.put("hesam", 1L); + + org.junit.Assert.assertEquals(expected, actual); + } + + @Test + void givenTextWithDuplicateWords_whenCallingCountWordOnGateway_thenWordCountReturnedAsMap() { + Map actual = testGateway.countWords("the dog chased the rabbit into the jungle"); + Map expected = new HashMap<>(); + expected.put("the", 3L); + expected.put("dog", 1L); + expected.put("chased", 1L); + expected.put("rabbit", 1L); + expected.put("into", 1L); + expected.put("jungle", 1L); + + org.junit.Assert.assertEquals(expected, actual); + } + +}