diff --git a/spring-boot-modules/spring-boot-libraries-3/pom.xml b/spring-boot-modules/spring-boot-libraries-3/pom.xml
index d0c1d345c0..988ce0bafe 100644
--- a/spring-boot-modules/spring-boot-libraries-3/pom.xml
+++ b/spring-boot-modules/spring-boot-libraries-3/pom.xml
@@ -6,9 +6,10 @@
spring-boot-libraries-3
- spring-boot-modules
- com.baeldung.spring-boot-modules
- 1.0.0-SNAPSHOT
+ com.baeldung
+ parent-boot-3
+ 0.0.1-SNAPSHOT
+ ../../parent-boot-3
@@ -16,12 +17,17 @@
org.springframework.boot
spring-boot-starter-data-jpa
+
org.springframework.kafka
spring-kafka
- ${spring-kafka.version}
+
+ org.postgresql
+ postgresql
+ ${postgresql.version}
+
org.springframework.modulith
spring-modulith-events-api
@@ -33,17 +39,11 @@
${spring-modulith-events-kafka.version}
- com.fasterxml.jackson.core
- jackson-databind
-
-
- com.fasterxml.jackson.core
- jackson-core
-
-
- com.fasterxml.jackson.core
- jackson-annotations
+ org.springframework.modulith
+ spring-modulith-starter-jpa
+ ${spring-modulith-events-kafka.version}
+
org.springframework.boot
spring-boot-starter-test
@@ -75,8 +75,10 @@
- com.h2database
- h2
+ org.testcontainers
+ postgresql
+ ${testcontainers.version}
+ test
@@ -90,10 +92,12 @@
17
- 1.1.3
- 1.19.6
+ 3.1.5
+ 1.1.2
+ 1.19.3
4.2.0
- 3.1.2
+ 42.3.1
-
\ No newline at end of file
+
+
diff --git a/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/Baeldung.java b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/Baeldung.java
index 1d309a8653..4b861a49c8 100644
--- a/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/Baeldung.java
+++ b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/Baeldung.java
@@ -6,32 +6,29 @@ import org.springframework.transaction.annotation.Transactional;
@Service
public class Baeldung {
+ private final ApplicationEventPublisher applicationEvents;
+ private final ArticleRepository articleRepository;
- private final ApplicationEventPublisher applicationEvents;
- private final ArticleRepository articleRepository;
+ public Baeldung(ApplicationEventPublisher applicationEvents, ArticleRepository articleRepository) {
+ this.applicationEvents = applicationEvents;
+ this.articleRepository = articleRepository;
+ }
+ @Transactional
+ public void createArticle(Article article) {
+ // ... business logic
+ validateArticle(article);
+ article = addArticleTags(article);
+ article = articleRepository.save(article);
- public Baeldung(ApplicationEventPublisher applicationEvents, ArticleRepository articleRepository) {
- this.applicationEvents = applicationEvents;
- this.articleRepository = articleRepository;
- }
+ applicationEvents.publishEvent(new ArticlePublishedEvent(article.slug(), article.title()));
+ }
- @Transactional
- public void createArticle(Article article) {
- // ... business logic
- validateArticle(article);
- article = addArticleTags(article);
- article = articleRepository.save(article);
+ private Article addArticleTags(Article article) {
+ return article;
+ }
- applicationEvents.publishEvent(new ArticlePublishedEvent(article.slug(), article.title()));
- }
-
-
- private Article addArticleTags(Article article) {
- return article;
- }
-
- private void validateArticle(Article article) {
- }
+ private void validateArticle(Article article) {
+ }
}
diff --git a/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/EventExternalizationConfig.java b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/EventExternalizationConfig.java
index 9564696d92..6555694df9 100644
--- a/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/EventExternalizationConfig.java
+++ b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/EventExternalizationConfig.java
@@ -4,12 +4,14 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
+import org.springframework.kafka.core.KafkaOperations;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
-import org.springframework.kafka.core.KafkaOperations;
import org.springframework.modulith.events.EventExternalizationConfiguration;
import org.springframework.modulith.events.RoutingTarget;
+import java.util.Objects;
+
@Configuration
class EventExternalizationConfig {
@@ -23,19 +25,30 @@ class EventExternalizationConfig {
)
.mapping(
ArticlePublishedEvent.class,
- it -> new ArticlePublishedKafkaEvent(it.slug(), it.title())
+ it -> new PostPublishedKafkaEvent(it.slug(), it.title())
+ )
+ .route(
+ WeeklySummaryPublishedEvent.class,
+ it -> RoutingTarget.forTarget("baeldung.articles.published").andKey(it.handle())
+ )
+ .mapping(
+ WeeklySummaryPublishedEvent.class,
+ it -> new PostPublishedKafkaEvent(it.handle(), it.heading())
)
.build();
}
- record ArticlePublishedKafkaEvent(String slug, String title) {
- }
-
-
@Bean
KafkaOperations kafkaOperations(KafkaProperties kafkaProperties) {
ProducerFactory producerFactory = new DefaultKafkaProducerFactory<>(kafkaProperties.buildProducerProperties());
return new KafkaTemplate<>(producerFactory);
}
+
+ record PostPublishedKafkaEvent(String slug, String title) {
+ PostPublishedKafkaEvent {
+ Objects.requireNonNull(slug, "Article Slug must not be null!");
+ }
+ }
+
}
diff --git a/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/WeeklySummaryPublishedEvent.java b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/WeeklySummaryPublishedEvent.java
new file mode 100644
index 0000000000..2ae8713099
--- /dev/null
+++ b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/WeeklySummaryPublishedEvent.java
@@ -0,0 +1,7 @@
+package com.baeldung.springmodulith.events.externalization;
+
+import org.springframework.modulith.events.Externalized;
+
+@Externalized
+record WeeklySummaryPublishedEvent(String handle, String heading) {
+}
\ No newline at end of file
diff --git a/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/infra/PublicationEvents.java b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/infra/PublicationEvents.java
new file mode 100644
index 0000000000..bf0b96e78b
--- /dev/null
+++ b/spring-boot-modules/spring-boot-libraries-3/src/main/java/com/baeldung/springmodulith/events/externalization/infra/PublicationEvents.java
@@ -0,0 +1,38 @@
+package com.baeldung.springmodulith.events.externalization.infra;
+
+import com.baeldung.springmodulith.events.externalization.ArticlePublishedEvent;
+import org.springframework.modulith.events.CompletedEventPublications;
+import org.springframework.modulith.events.IncompleteEventPublications;
+import org.springframework.stereotype.Component;
+
+import java.time.Duration;
+import java.time.Instant;
+
+@Component
+class PublicationEvents {
+ private final IncompleteEventPublications incompleteEvent;
+ private final CompletedEventPublications completeEvents;
+
+ public PublicationEvents(IncompleteEventPublications incompleteEvent, CompletedEventPublications completeEvents) {
+ this.incompleteEvent = incompleteEvent;
+ this.completeEvents = completeEvents;
+ }
+
+ public void resubmitUnpublishedEvents() {
+ incompleteEvent.resubmitIncompletePublicationsOlderThan(Duration.ofSeconds(60));
+
+ // or
+ incompleteEvent.resubmitIncompletePublications(it ->
+ it.getPublicationDate().isBefore(Instant.now().minusSeconds(60))
+ && it.getEvent() instanceof ArticlePublishedEvent);
+ }
+
+ public void clearPublishedEvents() {
+ completeEvents.deletePublicationsOlderThan(Duration.ofSeconds(60));
+
+ // or
+ completeEvents.deletePublications(it ->
+ it.getPublicationDate().isBefore(Instant.now().minusSeconds(60))
+ && it.getEvent() instanceof ArticlePublishedEvent);
+ }
+}
diff --git a/spring-boot-modules/spring-boot-libraries-3/src/main/resources/application.yml b/spring-boot-modules/spring-boot-libraries-3/src/main/resources/application.yml
index ad8f7ab4e6..c6797b57d0 100644
--- a/spring-boot-modules/spring-boot-libraries-3/src/main/resources/application.yml
+++ b/spring-boot-modules/spring-boot-libraries-3/src/main/resources/application.yml
@@ -1,4 +1,3 @@
-logging.level.org.springframework.orm.jpa: TRACE
spring.kafka:
bootstrap-servers: localhost:9092
@@ -10,3 +9,20 @@ spring.kafka:
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
auto-offset-reset: earliest
+
+spring.modulith:
+ republish-outstanding-events-on-restart: true
+ events.jdbc.schema-initialization.enabled: true
+
+logging.level.org.springframework.orm.jpa: TRACE
+
+spring:
+ datasource:
+ username: test_user
+ password: test_pass
+ jpa:
+ properties:
+ hibernate:
+ dialect: org.hibernate.dialect.PostgreSQLDialect
+ hbm2ddl.auto: create
+
diff --git a/spring-boot-modules/spring-boot-libraries-3/src/test/java/com/baeldung/springmodulith/events/externalization/EventsExternalizationLiveTest.java b/spring-boot-modules/spring-boot-libraries-3/src/test/java/com/baeldung/springmodulith/events/externalization/EventsExternalizationLiveTest.java
index 7f4f7a8224..a1b3dfe170 100644
--- a/spring-boot-modules/spring-boot-libraries-3/src/test/java/com/baeldung/springmodulith/events/externalization/EventsExternalizationLiveTest.java
+++ b/spring-boot-modules/spring-boot-libraries-3/src/test/java/com/baeldung/springmodulith/events/externalization/EventsExternalizationLiveTest.java
@@ -1,10 +1,8 @@
package com.baeldung.springmodulith.events.externalization;
-import static java.time.Duration.ofMillis;
-import static java.time.Duration.ofSeconds;
-import static org.assertj.core.api.Assertions.assertThat;
-import static org.testcontainers.shaded.org.awaitility.Awaitility.await;
-
+import com.baeldung.springmodulith.Application;
+import com.baeldung.springmodulith.events.externalization.listener.TestKafkaListenerConfig;
+import com.baeldung.springmodulith.events.externalization.listener.TestListener;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
@@ -12,14 +10,16 @@ import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.containers.KafkaContainer;
+import org.testcontainers.containers.PostgreSQLContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.shaded.org.awaitility.Awaitility;
import org.testcontainers.utility.DockerImageName;
-import com.baeldung.springmodulith.Application;
-import com.baeldung.springmodulith.events.externalization.listener.TestKafkaListenerConfig;
-import com.baeldung.springmodulith.events.externalization.listener.TestListener;
+import static java.time.Duration.ofMillis;
+import static java.time.Duration.ofSeconds;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.testcontainers.shaded.org.awaitility.Awaitility.await;
@Testcontainers
@SpringBootTest(classes = { Application.class, TestKafkaListenerConfig.class })
@@ -35,13 +35,20 @@ class EventsExternalizationLiveTest {
@Container
static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest"));
+ @Container
+ public static PostgreSQLContainer postgresqlContainer = new PostgreSQLContainer()
+ .withDatabaseName("test_db")
+ .withUsername("test_user")
+ .withPassword("test_pass");
+
@DynamicPropertySource
static void dynamicProperties(DynamicPropertyRegistry registry) {
registry.add("spring.kafka.bootstrap-servers", kafkaContainer::getBootstrapServers);
+ registry.add("spring.datasource.url", postgresqlContainer::getJdbcUrl);
}
static {
- Awaitility.setDefaultTimeout(ofSeconds(3));
+ Awaitility.setDefaultTimeout(ofSeconds(50));
Awaitility.setDefaultPollDelay(ofMillis(100));
}
@@ -86,4 +93,4 @@ class EventsExternalizationLiveTest {
.extracting(Article::title, Article::author)
.containsExactly("Introduction to Spring Boot", "John Doe");
}
-}
\ No newline at end of file
+}