From 0b49f4723089973d0c9cdc60473f7d2ea60e4a2e Mon Sep 17 00:00:00 2001 From: Christoph Strobl Date: Tue, 7 May 2019 10:30:18 +0200 Subject: [PATCH] DATAMONGO-2265 - Polishing. Fix count operation inside transaction and avoid superfluous client session instantiation. Default MongoDatabase emission in case of non active transaction, update documentation, move test to another package. Delay reactive collection re/creation in test to cope with issues in server version 4.1.10. Original Pull Request: #745 --- .../mongodb/ReactiveMongoDatabaseFactory.java | 2 - .../mongodb/ReactiveMongoDatabaseUtils.java | 59 ++- .../mongodb/ReactiveMongoResourceHolder.java | 26 +- .../ReactiveMongoTransactionManager.java | 18 +- .../mongodb/core/ReactiveMongoTemplate.java | 52 +-- .../ReactiveMongoDatabaseUtilsUnitTests.java | 8 +- ...ctiveMongoTransactionManagerUnitTests.java | 30 +- .../ReactiveTransactionIntegrationTests.java | 344 ++++++++++++++++++ .../ReactiveTransactionIntegrationTests.java | 128 ------- .../mongodb/test/util/MongoTestUtils.java | 23 +- src/main/asciidoc/new-features.adoc | 1 + .../client-session-transactions.adoc | 19 +- 12 files changed, 473 insertions(+), 237 deletions(-) create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java delete mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/rxtx/ReactiveTransactionIntegrationTests.java diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseFactory.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseFactory.java index 17b0f7d92..ac5032891 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseFactory.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseFactory.java @@ -13,13 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ - package org.springframework.data.mongodb; import reactor.core.publisher.Mono; import org.bson.codecs.configuration.CodecRegistry; - import org.springframework.dao.DataAccessException; import org.springframework.dao.support.PersistenceExceptionTranslator; import org.springframework.data.mongodb.core.MongoExceptionTranslator; diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtils.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtils.java index 60035bde5..e55b741cd 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtils.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtils.java @@ -1,11 +1,11 @@ /* - * Copyright 2018 the original author or authors. + * Copyright 2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * https://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -33,13 +33,14 @@ import com.mongodb.reactivestreams.client.MongoCollection; import com.mongodb.reactivestreams.client.MongoDatabase; /** - * Helper class for managing a {@link MongoDatabase} instances via {@link ReactiveMongoDatabaseFactory}. Used for + * Helper class for managing reactive {@link MongoDatabase} instances via {@link ReactiveMongoDatabaseFactory}. Used for * obtaining {@link ClientSession session bound} resources, such as {@link MongoDatabase} and {@link MongoCollection} * suitable for transactional usage. *

* Note: Intended for internal usage only. * * @author Mark Paluch + * @author Christoph Strobl * @since 2.2 */ public class ReactiveMongoDatabaseUtils { @@ -53,7 +54,7 @@ public class ReactiveMongoDatabaseUtils { * {@link com.mongodb.reactivestreams.client.ClientSession#hasActiveTransaction() active transaction}. * * @param databaseFactory the resource to check transactions for. Must not be {@literal null}. - * @return {@literal true} if the factory has an ongoing transaction. + * @return a {@link Mono} emitting {@literal true} if the factory has an ongoing transaction. */ public static Mono isTransactionActive(ReactiveMongoDatabaseFactory databaseFactory) { @@ -61,11 +62,13 @@ public class ReactiveMongoDatabaseUtils { return Mono.just(true); } - return TransactionSynchronizationManager.currentTransaction().map(it -> { + return TransactionSynchronizationManager.currentTransaction() // + .map(it -> { - ReactiveMongoResourceHolder holder = (ReactiveMongoResourceHolder) it.getResource(databaseFactory); - return holder != null && holder.hasActiveTransaction(); - }).onErrorResume(NoTransactionException.class, e -> Mono.just(false)); + ReactiveMongoResourceHolder holder = (ReactiveMongoResourceHolder) it.getResource(databaseFactory); + return holder != null && holder.hasActiveTransaction(); + }) // + .onErrorResume(NoTransactionException.class, e -> Mono.just(false)); } /** @@ -132,23 +135,23 @@ public class ReactiveMongoDatabaseUtils { private static Mono doGetMongoDatabase(@Nullable String dbName, ReactiveMongoDatabaseFactory factory, SessionSynchronization sessionSynchronization) { - Assert.notNull(factory, "Factory must not be null!"); + Assert.notNull(factory, "DatabaseFactory must not be null!"); return TransactionSynchronizationManager.currentTransaction() - .filter(TransactionSynchronizationManager::isSynchronizationActive).flatMap(synchronizationManager -> { + .filter(TransactionSynchronizationManager::isSynchronizationActive) // + .flatMap(synchronizationManager -> { - Mono session = doGetSession(synchronizationManager, factory, sessionSynchronization); + return doGetSession(synchronizationManager, factory, sessionSynchronization) // + .map(it -> getMongoDatabaseOrDefault(dbName, factory.withSession(it))); + }) + .onErrorResume(NoTransactionException.class, + e -> Mono.fromSupplier(() -> getMongoDatabaseOrDefault(dbName, factory))) + .defaultIfEmpty(getMongoDatabaseOrDefault(dbName, factory)); + } - return session.map(it -> { - - ReactiveMongoDatabaseFactory factoryToUse = factory.withSession(it); - return StringUtils.hasText(dbName) ? factoryToUse.getMongoDatabase(dbName) - : factoryToUse.getMongoDatabase(); - }); - - }).onErrorResume(NoTransactionException.class, e -> Mono.fromSupplier(() -> { - return StringUtils.hasText(dbName) ? factory.getMongoDatabase(dbName) : factory.getMongoDatabase(); - })); + private static MongoDatabase getMongoDatabaseOrDefault(@Nullable String dbName, + ReactiveMongoDatabaseFactory factory) { + return StringUtils.hasText(dbName) ? factory.getMongoDatabase(dbName) : factory.getMongoDatabase(); } private static Mono doGetSession(TransactionSynchronizationManager synchronizationManager, @@ -161,14 +164,8 @@ public class ReactiveMongoDatabaseUtils { if (registeredHolder != null && (registeredHolder.hasSession() || registeredHolder.isSynchronizedWithTransaction())) { - return createClientSession(dbFactory).map(session -> { - - if (!registeredHolder.hasSession()) { - registeredHolder.setSession(session); - } - - return registeredHolder.getSession(); - }); + return registeredHolder.hasSession() ? Mono.just(registeredHolder.getSession()) + : createClientSession(dbFactory).map(registeredHolder::setSessionIfAbsent); } if (SessionSynchronization.ON_ACTUAL_TRANSACTION.equals(sessionSynchronization)) { @@ -246,7 +243,9 @@ public class ReactiveMongoDatabaseUtils { return Mono.defer(() -> { if (status == TransactionSynchronization.STATUS_ROLLED_BACK && isTransactionActive(this.resourceHolder)) { - return Mono.from(resourceHolder.getRequiredSession().abortTransaction()).then(super.afterCompletion(status)); + + return Mono.from(resourceHolder.getRequiredSession().abortTransaction()) // + .then(super.afterCompletion(status)); } return super.afterCompletion(status); diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoResourceHolder.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoResourceHolder.java index 49bc77570..cb19f15ba 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoResourceHolder.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoResourceHolder.java @@ -1,11 +1,11 @@ /* - * Copyright 2018 the original author or authors. + * Copyright 2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * https://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -22,8 +22,8 @@ import org.springframework.transaction.support.ResourceHolderSupport; import com.mongodb.reactivestreams.client.ClientSession; /** - * MongoDB specific resource holder, wrapping a {@link ClientSession}. {@link MongoTransactionManager} binds instances - * of this class to the subscriber context. + * MongoDB specific resource holder, wrapping a {@link ClientSession}. {@link ReactiveMongoTransactionManager} binds + * instances of this class to the subscriber context. *

* Note: Intended for internal usage only. * @@ -95,6 +95,24 @@ class ReactiveMongoResourceHolder extends ResourceHolderSupport { return session != null; } + /** + * If the {@link ReactiveMongoResourceHolder} is {@link #hasSession() not already associated} with a + * {@link ClientSession} the given value is {@link #setSession(ClientSession) set} and returned, otherwise the current + * bound session is returned. + * + * @param session + * @return + */ + @Nullable + public ClientSession setSessionIfAbsent(@Nullable ClientSession session) { + + if (!hasSession()) { + setSession(session); + } + + return session; + } + /** * @return {@literal true} if the session is active and has not been closed. */ diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java index 3da00bb48..9bd0fbfe0 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java @@ -1,11 +1,11 @@ /* - * Copyright 2018 the original author or authors. + * Copyright 2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * https://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -104,7 +104,7 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction public ReactiveMongoTransactionManager(ReactiveMongoDatabaseFactory databaseFactory, @Nullable TransactionOptions options) { - Assert.notNull(databaseFactory, "DbFactory must not be null!"); + Assert.notNull(databaseFactory, "DatabaseFactory must not be null!"); this.databaseFactory = databaseFactory; this.options = options; @@ -164,9 +164,10 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction logger.debug(String.format("Started transaction for session %s.", debugString(resourceHolder.getSession()))); } - }).onErrorMap( - ex -> new TransactionSystemException(String.format("Could not start Mongo transaction for session %s.", - debugString(mongoTransactionObject.getSession())), ex)) + })// + .onErrorMap( + ex -> new TransactionSystemException(String.format("Could not start Mongo transaction for session %s.", + debugString(mongoTransactionObject.getSession())), ex)) .doOnSuccess(resourceHolder -> { synchronizationManager.bindResource(getRequiredDatabaseFactory(), resourceHolder); @@ -316,7 +317,7 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction */ public void setDatabaseFactory(ReactiveMongoDatabaseFactory databaseFactory) { - Assert.notNull(databaseFactory, "DbFactory must not be null!"); + Assert.notNull(databaseFactory, "DatabaseFactory must not be null!"); this.databaseFactory = databaseFactory; } @@ -362,7 +363,7 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction private ReactiveMongoDatabaseFactory getRequiredDatabaseFactory() { Assert.state(databaseFactory != null, - "MongoTransactionManager operates upon a ReactiveMongoDatabaseFactory. Did you forget to provide one? It's required."); + "ReactiveMongoTransactionManager operates upon a ReactiveMongoDatabaseFactory. Did you forget to provide one? It's required."); return databaseFactory; } @@ -525,6 +526,5 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction public void flush() { throw new UnsupportedOperationException("flush() not supported"); } - } } diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java index d0c0592b5..944e5952d 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java @@ -707,17 +707,6 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati } } - /* - * (non-Javadoc) - * @see org.springframework.data.mongodb.core.ReactiveMongoOperations#getCollection(java.lang.String) - */ - public Mono> getCollection2(final String collectionName) { - - Assert.notNull(collectionName, "Collection name must not be null!"); - - return doGetDatabase().map(it -> it.getCollection(collectionName)); - } - /* * (non-Javadoc) * @see org.springframework.data.mongodb.core.ReactiveMongoOperations#collectionExists(java.lang.Class) @@ -1247,10 +1236,25 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati LOGGER.debug("Executing count: {} in collection: {}", serializeToJsonSafely(filter), collectionName); } - return collection.count(filter, options); + return doCount(collectionName, filter, options); }); } + /** + * Run the actual count operation against the collection with given name. + * + * @param collectionName the name of the collection to count matching documents in. + * @param filter the filter to apply. Must not be {@literal null}. + * @param options options to apply. Like collation and the such. + * @return + */ + protected Mono doCount(String collectionName, Document filter, CountOptions options) { + + return ReactiveMongoDatabaseUtils.isTransactionActive(mongoDatabaseFactory) // + .flatMap(txActive -> createMono(collectionName, + collection -> txActive ? collection.countDocuments(filter, options) : collection.count(filter, options))); + } + /* * (non-Javadoc) * @see org.springframework.data.mongodb.core.ReactiveMongoOperations#insert(reactor.core.publisher.Mono) @@ -3218,32 +3222,16 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati /* * (non-Javadoc) - * @see org.springframework.data.mongodb.core.ReactiveMongoTemplate#count(org.springframework.data.mongodb.core.query.Query, java.lang.Class, java.lang.String) + * @see org.springframework.data.mongodb.core.ReactiveMongoTemplate#count(java.lang.String, org.bson.Document, com.mongodb.client.model.CountOptions) */ @Override - public Mono count(Query query, @Nullable Class entityClass, String collectionName) { + public Mono doCount(String collectionName, Document filter, CountOptions options) { if (!session.hasActiveTransaction()) { - return super.count(query, entityClass, collectionName); + return super.doCount(collectionName, filter, options); } - return createMono(collectionName, collection -> { - - Document filter = query == null ? null - : delegate.queryMapper.getMappedObject(query.getQueryObject(), - entityClass == null ? null : delegate.mappingContext.getPersistentEntity(entityClass)); - - CountOptions options = new CountOptions(); - if (query != null) { - query.getCollation().map(Collation::toMongoCollation).ifPresent(options::collation); - } - - if (LOGGER.isDebugEnabled()) { - LOGGER.debug("Executing count: {} in collection: {}", serializeToJsonSafely(filter), collectionName); - } - - return collection.countDocuments(filter, options); - }); + return createMono(collectionName, collection -> collection.countDocuments(filter, options)); } } diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtilsUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtilsUnitTests.java index c23fb17bf..ec840a4c0 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtilsUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoDatabaseUtilsUnitTests.java @@ -16,7 +16,7 @@ package org.springframework.data.mongodb; import static org.assertj.core.api.Assertions.*; -import static org.mockito.ArgumentMatchers.*; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.*; import reactor.core.publisher.Mono; @@ -27,7 +27,6 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; - import org.springframework.transaction.reactive.TransactionSynchronizationManager; import org.springframework.transaction.reactive.TransactionalOperator; import org.springframework.transaction.support.DefaultTransactionDefinition; @@ -40,6 +39,7 @@ import com.mongodb.session.ServerSession; * Unit tests for {@link ReactiveMongoDatabaseUtils}. * * @author Mark Paluch + * @author Christoph Strobl */ @RunWith(MockitoJUnitRunner.class) public class ReactiveMongoDatabaseUtilsUnitTests { @@ -128,8 +128,6 @@ public class ReactiveMongoDatabaseUtilsUnitTests { verify(session).startTransaction(); verify(session).abortTransaction(); - - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); + verify(session).close(); } } diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoTransactionManagerUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoTransactionManagerUnitTests.java index a5d16933a..b8f2ec152 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoTransactionManagerUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveMongoTransactionManagerUnitTests.java @@ -27,7 +27,6 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; - import org.springframework.data.mongodb.core.ReactiveMongoTemplate; import org.springframework.transaction.TransactionDefinition; import org.springframework.transaction.reactive.TransactionalOperator; @@ -42,6 +41,7 @@ import com.mongodb.session.ServerSession; * Unit tests for {@link ReactiveMongoTransactionManager}. * * @author Mark Paluch + * @author Christoph Strobl */ @RunWith(MockitoJUnitRunner.class) public class ReactiveMongoTransactionManagerUnitTests { @@ -64,6 +64,9 @@ public class ReactiveMongoTransactionManagerUnitTests { when(databaseFactory.getMongoDatabase()).thenReturn(db); when(databaseFactory2.getMongoDatabase()).thenReturn(db2); + + when(session.getServerSession()).thenReturn(serverSession); + when(session2.getServerSession()).thenReturn(serverSession); } @After @@ -95,8 +98,7 @@ public class ReactiveMongoTransactionManagerUnitTests { verify(session).startTransaction(); verify(session).commitTransaction(); - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); + verify(session).close(); } @Test // DATAMONGO-2265 @@ -124,9 +126,7 @@ public class ReactiveMongoTransactionManagerUnitTests { verify(session).startTransaction(); verify(session).commitTransaction(); - - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); + verify(session).close(); } @Test // DATAMONGO-2265 @@ -151,9 +151,7 @@ public class ReactiveMongoTransactionManagerUnitTests { verify(session).startTransaction(); verify(session).abortTransaction(); - - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); + verify(session).close(); } @Test // DATAMONGO-2265 @@ -190,11 +188,8 @@ public class ReactiveMongoTransactionManagerUnitTests { verify(databaseFactory, times(1)).withSession(eq(session)); verify(databaseFactory, never()).withSession(eq(session2)); - // Bug in TransactionalOperator, should be 2 - verify(db, times(1)).drop(); + verify(db, times(2)).drop(); - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); verify(session2, never()).close(); } @@ -236,9 +231,8 @@ public class ReactiveMongoTransactionManagerUnitTests { verify(db).drop(); verify(db2).drop(); - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); - // verify(session2).close(); + verify(session).close(); + verify(session2).close(); } @Test // DATAMONGO-2265 @@ -264,8 +258,6 @@ public class ReactiveMongoTransactionManagerUnitTests { verify(session).startTransaction(); verify(session).commitTransaction(); - - // TODO: Bug in doCleanupAfterCompletion - // verify(session).close(); + verify(session).close(); } } diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java new file mode 100644 index 000000000..24e1d84da --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java @@ -0,0 +1,344 @@ +/* + * Copyright 2019 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.RequiredArgsConstructor; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.test.StepVerifier; + +import java.time.Duration; + +import org.bson.types.ObjectId; +import org.junit.AfterClass; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.ClassRule; +import org.junit.Test; +import org.junit.rules.RuleChain; +import org.junit.runner.RunWith; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.support.GenericApplicationContext; +import org.springframework.data.mongodb.config.AbstractReactiveMongoConfiguration; +import org.springframework.data.mongodb.core.ReactiveMongoOperations; +import org.springframework.data.mongodb.core.mapping.Document; +import org.springframework.data.mongodb.core.query.Query; +import org.springframework.data.mongodb.test.util.MongoTestUtils; +import org.springframework.data.mongodb.test.util.MongoVersionRule; +import org.springframework.data.mongodb.test.util.ReplicaSet; +import org.springframework.data.util.Version; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.reactive.TransactionalOperator; +import org.springframework.transaction.support.DefaultTransactionDefinition; + +import com.mongodb.reactivestreams.client.MongoClient; +import com.mongodb.reactivestreams.client.MongoClients; + +/** + * Integration tests for reactive transaction management. + * + * @author Mark Paluch + * @author Christoph Strobl + */ +@RunWith(SpringJUnit4ClassRunner.class) +public class ReactiveTransactionIntegrationTests { + + public static @ClassRule RuleChain TEST_RULES = RuleChain.outerRule(MongoVersionRule.atLeast(Version.parse("4.0.0"))) + .around(ReplicaSet.required()); + + private static final String DATABASE = "rxtx-test"; + PersonService personService; + ReactiveMongoOperations operations; + static GenericApplicationContext context; + + @BeforeClass + public static void init() { + context = new AnnotationConfigApplicationContext(TestMongoConfig.class, PersonService.class); + } + + @AfterClass + public static void after() { + context.close(); + } + + @Before + public void setUp() { + + personService = context.getBean(PersonService.class); + operations = context.getBean(ReactiveMongoOperations.class); + + try (MongoClient client = MongoClients.create()) { + + Flux.merge( // + MongoTestUtils.createOrReplaceCollection(DATABASE, operations.getCollectionName(Person.class), client), + MongoTestUtils.createOrReplaceCollection(DATABASE, operations.getCollectionName(EventLog.class), client) // + ).then().as(StepVerifier::create).verifyComplete(); + } + } + + @Test // DATAMONGO-2265 + public void shouldRollbackAfterException() { + + personService.savePersonErrors(new Person(null, "Walter", "White")) // + .as(StepVerifier::create) // + .verifyError(RuntimeException.class); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(0L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void shouldRollbackAfterExceptionOfTxAnnotatedMethod() { + + personService.declarativeSavePersonErrors(new Person(null, "Walter", "White")) // + .as(StepVerifier::create) // + .verifyError(RuntimeException.class); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(0L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void commitShouldPersistTxEntries() { + + personService.savePerson(new Person(null, "Walter", "White")) // + .as(StepVerifier::create) // + .expectNextCount(1) // + .verifyComplete(); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(1L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void commitShouldPersistTxEntriesOfTxAnnotatedMethod() { + + personService.declarativeSavePerson(new Person(null, "Walter", "White")) // + .as(StepVerifier::create) // + .expectNextCount(1) // + .verifyComplete(); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(1L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void commitShouldPersistTxEntriesAcrossCollections() { + + personService.saveWithLogs(new Person(null, "Walter", "White")) // + .then() // + .as(StepVerifier::create) // + .verifyComplete(); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(1L) // + .verifyComplete(); + + operations.count(new Query(), EventLog.class) // + .as(StepVerifier::create) // + .expectNext(4L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void rollbackShouldAbortAcrossCollections() { + + personService.saveWithErrorLogs(new Person(null, "Walter", "White")) // + .then() // + .as(StepVerifier::create) // + .verifyError(); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(0L) // + .verifyComplete(); + + operations.count(new Query(), EventLog.class) // + .as(StepVerifier::create) // + .expectNext(0L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void countShouldWorkInsideTransaction() { + + personService.countDuringTx(new Person(null, "Walter", "White")) // + .as(StepVerifier::create) // + .expectNext(1L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void emitMultipleElementsDuringTransaction() { + + personService.saveWithLogs(new Person(null, "Walter", "White")) // + .as(StepVerifier::create) // + .expectNextCount(4L) // + .verifyComplete(); + } + + @Test // DATAMONGO-2265 + public void errorAfterTxShouldNotAffectPreviousStep() { + + personService.savePerson(new Person(null, "Walter", "White")) // + .delayElement(Duration.ofMillis(10)) // + .then(Mono.error(new RuntimeException("my big bad evil error"))).as(StepVerifier::create) // + .expectError() // + .verify(); + + operations.count(new Query(), Person.class) // + .as(StepVerifier::create) // + .expectNext(1L) // + .verifyComplete(); + } + + @Configuration + static class TestMongoConfig extends AbstractReactiveMongoConfiguration { + + @Override + public MongoClient reactiveMongoClient() { + return MongoClients.create("mongodb://localhost"); + } + + @Override + protected String getDatabaseName() { + return DATABASE; + } + + @Bean + public ReactiveMongoTransactionManager transactionManager(ReactiveMongoDatabaseFactory factory) { + return new ReactiveMongoTransactionManager(factory); + } + } + + @RequiredArgsConstructor + static class PersonService { + + final ReactiveMongoOperations operations; + final ReactiveMongoTransactionManager manager; + + public Mono savePersonErrors(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return operations.save(person) // + . flatMap(it -> Mono.error(new RuntimeException("poof!"))) // + .as(transactionalOperator::transactional); + } + + public Mono savePerson(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return operations.save(person) // + .flatMap(Mono::just) // + .as(transactionalOperator::transactional); + } + + public Mono countDuringTx(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return operations.save(person) // + .then(operations.count(new Query(), Person.class)) // + .as(transactionalOperator::transactional); + } + + public Flux saveWithLogs(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return Flux.merge(operations.save(new EventLog(new ObjectId(), "beforeConvert")), // + operations.save(new EventLog(new ObjectId(), "afterConvert")), // + operations.save(new EventLog(new ObjectId(), "beforeInsert")), // + operations.save(person), // + operations.save(new EventLog(new ObjectId(), "afterInsert"))) // + .thenMany(operations.query(EventLog.class).all()) // + .as(transactionalOperator::transactional); + } + + public Flux saveWithErrorLogs(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return Flux.merge(operations.save(new EventLog(new ObjectId(), "beforeConvert")), // + operations.save(new EventLog(new ObjectId(), "afterConvert")), // + operations.save(new EventLog(new ObjectId(), "beforeInsert")), // + operations.save(person), // + operations.save(new EventLog(new ObjectId(), "afterInsert"))) // + . flatMap(it -> Mono.error(new RuntimeException("poof!"))) // + .as(transactionalOperator::transactional); + } + + @Transactional + public Mono declarativeSavePerson(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return operations.save(person) // + .flatMap(Mono::just) // + .as(transactionalOperator::transactional); + } + + @Transactional + public Mono declarativeSavePersonErrors(Person person) { + + TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, + new DefaultTransactionDefinition()); + + return operations.save(person) // + . flatMap(it -> Mono.error(new RuntimeException("poof!"))) // + .as(transactionalOperator::transactional); + } + } + + @Data + @AllArgsConstructor + @Document("person-rx") + static class Person { + + ObjectId id; + String firstname, lastname; + } + + @Data + @AllArgsConstructor + static class EventLog { + + ObjectId id; + String action; + } +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/rxtx/ReactiveTransactionIntegrationTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/rxtx/ReactiveTransactionIntegrationTests.java deleted file mode 100644 index e192a03a2..000000000 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/rxtx/ReactiveTransactionIntegrationTests.java +++ /dev/null @@ -1,128 +0,0 @@ -/* - * Copyright 2018 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.data.mongodb.rxtx; - -import lombok.AllArgsConstructor; -import lombok.Data; -import lombok.RequiredArgsConstructor; -import reactor.core.publisher.Mono; -import reactor.test.StepVerifier; - -import org.bson.types.ObjectId; -import org.junit.Test; - -import org.springframework.context.annotation.AnnotationConfigApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; -import org.springframework.data.mongodb.ReactiveMongoDatabaseFactory; -import org.springframework.data.mongodb.ReactiveMongoTransactionManager; -import org.springframework.data.mongodb.config.AbstractReactiveMongoConfiguration; -import org.springframework.data.mongodb.core.MongoTemplate; -import org.springframework.data.mongodb.core.ReactiveMongoOperations; -import org.springframework.data.mongodb.core.query.Query; -import org.springframework.stereotype.Service; -import org.springframework.transaction.reactive.TransactionalOperator; -import org.springframework.transaction.support.DefaultTransactionDefinition; - -import com.mongodb.reactivestreams.client.MongoClient; -import com.mongodb.reactivestreams.client.MongoClients; - -/** - * Integration tests for reactive transaction management. - * - * @author Mark Paluch - */ -public class ReactiveTransactionIntegrationTests { - - @Configuration - static class TestMongoConfig extends AbstractReactiveMongoConfiguration { - - @Override - public MongoClient reactiveMongoClient() { - return MongoClients.create("mongodb://localhost"); - } - - @Override - protected String getDatabaseName() { - return "test"; - } - - @Bean - public ReactiveMongoTransactionManager transactionManager(ReactiveMongoDatabaseFactory factory) { - return new ReactiveMongoTransactionManager(factory); - } - } - - @Service - @RequiredArgsConstructor - static class PersonService { - - final ReactiveMongoOperations operations; - final ReactiveMongoTransactionManager manager; - - public Mono savePerson(Person person) { - - TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, - new DefaultTransactionDefinition()); - - return operations.save(person). flatMap(it -> { - return Mono.error(new RuntimeException("poof!")); - }).as(transactionalOperator::transactional); - } - } - - @Test - public void shouldRollbackAfterException() { - - AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(TestMongoConfig.class, - PersonService.class); - - ReactiveMongoOperations operations = context.getBean(ReactiveMongoOperations.class); - - MongoTemplate template = new MongoTemplate(new com.mongodb.MongoClient("localhost"), "test"); - - template.dropCollection(Person.class); - template.dropCollection(EventLog.class); - - template.createCollection(Person.class); - template.createCollection(EventLog.class); - - PersonService personService = context.getBean(PersonService.class); - - personService.savePerson(new Person(null, "Walter", "White")) // - .as(StepVerifier::create) // - .verifyError(RuntimeException.class); - - operations.count(new Query(), Person.class).as(StepVerifier::create).expectNext(0L).verifyComplete(); - } - - @Data - @AllArgsConstructor - static class Person { - - ObjectId id; - String firstname, lastname; - } - - @Data - @AllArgsConstructor - static class EventLog { - - ObjectId id; - String action; - } - -} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java index 9c28be7e8..955a0ef9c 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/test/util/MongoTestUtils.java @@ -81,7 +81,8 @@ public class MongoTestUtils { .withWriteConcern(WriteConcern.MAJORITY).withReadPreference(ReadPreference.primary()); return Mono.from(database.getCollection(collectionName).drop()) // - .then(Mono.from(database.createCollection(collectionName))) // + .delayElement(Duration.ofMillis(10)) // server replication time + .then(Mono.from(database.createCollection(collectionName))) .delayElement(Duration.ofMillis(10)); // server replication time } @@ -102,6 +103,26 @@ public class MongoTestUtils { .verifyComplete(); } + /** + * Create a {@link com.mongodb.client.MongoCollection} if it does not exist, or drop and recreate it if it does and + * verify operation result. + * + * @param dbName must not be {@literal null}. + * @param collectionName must not be {@literal null}. + * @param client must not be {@literal null}. + */ + public static void dropCollectionNow(String dbName, String collectionName, + com.mongodb.reactivestreams.client.MongoClient client) { + + com.mongodb.reactivestreams.client.MongoDatabase database = client.getDatabase(dbName) + .withWriteConcern(WriteConcern.MAJORITY).withReadPreference(ReadPreference.primary()); + + Mono.from(database.getCollection(collectionName).drop()) // + .as(StepVerifier::create) // + .expectNext(Success.SUCCESS) // + .verifyComplete(); + } + /** * Create a new {@link com.mongodb.MongoClient} with defaults suitable for replica set usage. * diff --git a/src/main/asciidoc/new-features.adoc b/src/main/asciidoc/new-features.adoc index 36b94e4d3..3f8739da8 100644 --- a/src/main/asciidoc/new-features.adoc +++ b/src/main/asciidoc/new-features.adoc @@ -17,6 +17,7 @@ * <> from domain types. * SpEL support in for expressions in `@Indexed`. * Annotation-based Collation support through `@Document` and `@Query`. +* Declarative reactive transactions using <>. [[new-features.2-1-0]] == What's New in Spring Data MongoDB 2.1 diff --git a/src/main/asciidoc/reference/client-session-transactions.adoc b/src/main/asciidoc/reference/client-session-transactions.adoc index a79fcd235..a0584841b 100644 --- a/src/main/asciidoc/reference/client-session-transactions.adoc +++ b/src/main/asciidoc/reference/client-session-transactions.adoc @@ -138,6 +138,7 @@ template.withSession(session) The preceding example lets you have full control over transactional behavior while using the session scoped `MongoOperations` instance within the callback to ensure the session is passed on to every server call. To avoid some of the overhead that comes with this approach, you can use a `TransactionTemplate` to take away some of the noise of manual transaction flow. +[[mongo.transactions.transaction-template]] == Transactions with `TransactionTemplate` Spring Data MongoDB transactions support a `TransactionTemplate`. The following example shows how to create and use a `TransactionTemplate`: @@ -173,6 +174,7 @@ txTemplate.execute(new TransactionCallbackWithoutResult() { CAUTION: Changing state of `MongoTemplate` during runtime (as you might think would be possible in item 1 of the preceding listing) can cause threading and visibility issues. +[[mongo.transactions.tx-manager]] == Transactions with `MongoTransactionManager` `MongoTransactionManager` is the gateway to the well known Spring transaction support. It lets applications use https://docs.spring.io/spring/docs/{springVersion}/spring-framework-reference/html/transaction.html[the managed transaction features of Spring]. @@ -215,6 +217,7 @@ public class StateService { NOTE: `@Transactional(readOnly = true)` advises `MongoTransactionManager` to also start a transaction that adds the `ClientSession` to outgoing requests. +[[mongo.transactions.reactive]] == Reactive Transactions Same as with the reactive `ClientSession` support, the `ReactiveMongoTemplate` offers dedicated methods for operating @@ -254,6 +257,7 @@ Mono result = Mono The culprit of the above operation is in keeping the main flows `DeleteResult` instead of the transaction outcome published via either `commitTransaction()` or `abortTransaction()`, which leads to a rather complicated setup. +[[mongo.transactions.reactive-operator]] == Transactions with `TransactionalOperator` Spring Data MongoDB transactions support a `TransactionalOperator`. The following example shows how to create and use a `TransactionalOperator`: @@ -274,19 +278,20 @@ Step step = // ...; template.insert(step); Mono process(step) - .then(template.update(Step.class).apply(Update.set("state", …)) - .as(rxtx::transactional) <3> - .then(); + .then(template.update(Step.class).apply(Update.set("state", …)) + .as(rxtx::transactional) <3> + .then(); ---- <1> Enable transaction synchronization for Transactional participation. <2> Create the `TransactionalOperator` using the provided `ReactiveTransactionManager`. <3> `TransactionalOperator.transactional(…)` provides transaction management for all upstream operations. ==== +[[mongo.transactions.reactive-tx-manager]] == Transactions with `ReactiveMongoTransactionManager` `ReactiveMongoTransactionManager` is the gateway to the well known Spring transaction support. -It lets applications use https://docs.spring.io/spring/docs/{springVersion}/spring-framework-reference/html/transaction.html[the managed transaction features of Spring]. +It allows applications to leverage https://docs.spring.io/spring/docs/{springVersion}/spring-framework-reference/html/transaction.html[the managed transaction features of Spring]. The `ReactiveMongoTransactionManager` binds a `ClientSession` to the subscriber `Context`. `ReactiveMongoTemplate` detects the session and operates on these resources which are associated with the transaction accordingly. `ReactiveMongoTemplate` can also participate in other, ongoing transactions. @@ -307,15 +312,15 @@ static class Config extends AbstractMongoConfiguration { // ... } -@Component +@Service public class StateService { @Transactional Mono someBusinessFunction(Step step) { <2> return template.insert(step) - .then(process(step)) - .then(template.update(Step.class).apply(Update.set("state", …)); + .then(process(step)) + .then(template.update(Step.class).apply(Update.set("state", …)); }; });