From dc1148624bca2719e60f45067623da9c24d4ff07 Mon Sep 17 00:00:00 2001 From: Oliver Drotbohm Date: Mon, 21 Oct 2024 13:03:47 +0200 Subject: [PATCH] GH-806 - MongoDB archive mode now uses aggregation to mark event publications completed. --- .../MongoDbEventPublicationRepository.java | 56 ++++++++++++++----- .../src/main/java/example/Application.java | 10 +++- .../example/ApplicationIntegrationTests.java | 9 +++ 3 files changed, 60 insertions(+), 15 deletions(-) diff --git a/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java index 10b096b0..77577cb5 100644 --- a/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java +++ b/spring-modulith-events/spring-modulith-events-mongodb/src/main/java/org/springframework/modulith/events/mongodb/MongoDbEventPublicationRepository.java @@ -15,6 +15,7 @@ */ package org.springframework.modulith.events.mongodb; +import static org.springframework.data.mongodb.core.aggregation.Aggregation.*; import static org.springframework.data.mongodb.core.query.Criteria.*; import static org.springframework.data.mongodb.core.query.Query.*; @@ -24,8 +25,12 @@ import java.util.Objects; import java.util.Optional; import java.util.UUID; +import org.bson.Document; +import org.springframework.data.annotation.Id; import org.springframework.data.domain.Sort; import org.springframework.data.mongodb.core.MongoTemplate; +import org.springframework.data.mongodb.core.aggregation.Fields; +import org.springframework.data.mongodb.core.aggregation.MergeOperation.WhenDocumentsMatch; import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.data.mongodb.core.query.Query; import org.springframework.data.mongodb.core.query.Update; @@ -96,7 +101,8 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { @Override public void markCompleted(Object event, PublicationTargetIdentifier identifier, Instant completionDate) { - var query = byEventAndListenerId(event, identifier); + var criteria = byEventAndListenerId(event, identifier); + var query = defaultQuery(criteria); var update = Update.update(COMPLETION_DATE, completionDate); if (completionMode == CompletionMode.DELETE) { @@ -105,9 +111,8 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { } else if (completionMode == CompletionMode.ARCHIVE) { - mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class, collection); - var completedEvent = mongoTemplate.findAndRemove(query, MongoDbEventPublication.class, collection); - mongoTemplate.save(completedEvent, archiveCollection); + markCompleted(criteria, completionDate); + } else { mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class, collection); @@ -121,21 +126,20 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { @Override public void markCompleted(UUID identifier, Instant completionDate) { - var criteria = query(where(ID).is(identifier)); + var criteria = where(ID).is(identifier).and(COMPLETION_DATE).isNull(); + var query = query(criteria); var update = Update.update(COMPLETION_DATE, completionDate); if (completionMode == CompletionMode.DELETE) { - mongoTemplate.remove(criteria, MongoDbEventPublication.class, collection); + mongoTemplate.remove(query, MongoDbEventPublication.class, collection); } else if (completionMode == CompletionMode.ARCHIVE) { - mongoTemplate.findAndModify(criteria, update, MongoDbEventPublication.class, collection); - var completedEvent = mongoTemplate.findAndRemove(criteria, MongoDbEventPublication.class, collection); - mongoTemplate.save(completedEvent, archiveCollection); + markCompleted(criteria, completionDate); } else { - mongoTemplate.findAndModify(criteria, update, MongoDbEventPublication.class, collection); + mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class, collection); } } @@ -168,7 +172,7 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { public Optional findIncompletePublicationsByEventAndTargetIdentifier( Object event, PublicationTargetIdentifier targetIdentifier) { - var results = readMapped(byEventAndListenerId(event, targetIdentifier)); + var results = readMapped(defaultQuery(byEventAndListenerId(event, targetIdentifier))); // if there are several events with exactly the same payload we return the oldest one first return results.isEmpty() ? Optional.empty() : Optional.of(results.get(0)); @@ -230,13 +234,13 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { } - private Query byEventAndListenerId(Object event, PublicationTargetIdentifier identifier) { + private Criteria byEventAndListenerId(Object event, PublicationTargetIdentifier identifier) { var eventAsMongoType = mongoTemplate.getConverter().convertToMongoType(event, TypeInformation.OBJECT); - return defaultQuery(where(EVENT).is(eventAsMongoType) // + return where(EVENT).is(eventAsMongoType) // .and(LISTENER_ID).is(identifier.getValue()) - .and(COMPLETION_DATE).isNull()); + .and(COMPLETION_DATE).isNull(); } private static MongoDbEventPublication domainToDocument(TargetEventPublication publication) { @@ -256,6 +260,28 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { return query(criteria).with(DEFAULT_SORT); } + private void markCompleted(Criteria lookup, Instant now) { + + var aggregation = newAggregation(MongoDbEventPublication.class, + + match(lookup), + + addFields() + .addFieldWithValue(COMPLETION_DATE, now) + .build(), + + merge() + .intoCollection(archiveCollection) + .on(ID) + .whenMatched(WhenDocumentsMatch.keepExistingDocument()) + .build()); + + mongoTemplate + .aggregate(aggregation, collection, Document.class) + .forEach(it -> mongoTemplate.remove(query(where(Fields.UNDERSCORE_ID).is(it.get(Fields.UNDERSCORE_ID))), + collection)); + } + private static class MongoDbEventPublicationAdapter implements TargetEventPublication { private final MongoDbEventPublication publication; @@ -326,4 +352,6 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository { return Objects.hash(publication); } } + + record IdOnly(@Id UUID id) {} } diff --git a/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/main/java/example/Application.java b/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/main/java/example/Application.java index d48ed011..b2f381da 100644 --- a/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/main/java/example/Application.java +++ b/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/main/java/example/Application.java @@ -15,6 +15,9 @@ */ package example; +import example.order.OrderManagement; + +import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; /** @@ -25,5 +28,10 @@ import org.springframework.boot.autoconfigure.SpringBootApplication; @SpringBootApplication public class Application { - public static void main(String... args) {} + public static void main(String... args) { + + var context = SpringApplication.run(Application.class, args); + + context.getBean(OrderManagement.class).complete(); + } } diff --git a/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/test/java/example/ApplicationIntegrationTests.java b/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/test/java/example/ApplicationIntegrationTests.java index fab52da8..8d220d5f 100644 --- a/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/test/java/example/ApplicationIntegrationTests.java +++ b/spring-modulith-examples/spring-modulith-example-epr-mongodb/src/test/java/example/ApplicationIntegrationTests.java @@ -23,6 +23,7 @@ import java.util.Collection; import org.bson.UuidRepresentation; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.SpringApplication; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.TestConfiguration; import org.springframework.boot.testcontainers.service.connection.ServiceConnection; @@ -48,6 +49,14 @@ import com.mongodb.client.MongoClients; @Testcontainers(disabledWithoutDocker = true) class ApplicationIntegrationTests { + public static void main(String[] args) { + + SpringApplication.from(Application::main) + .with(MongoDbInfrastructureConfiguration.class) + .run(args) + .getApplicationContext(); + } + @TestConfiguration static class MongoDbInfrastructureConfiguration {