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
This commit is contained in:
Christoph Strobl
2019-05-07 10:30:18 +02:00
parent 5c10a5821b
commit 0b49f47230
12 changed files with 473 additions and 237 deletions

View File

@@ -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;

View File

@@ -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.
* <p />
* <strong>Note:</strong> 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<Boolean> 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<MongoDatabase> 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<ClientSession> 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<ClientSession> 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);

View File

@@ -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.
* <p />
* <strong>Note:</strong> 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.
*/

View File

@@ -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");
}
}
}

View File

@@ -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<MongoCollection<Document>> 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<Long> 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<Long> count(Query query, @Nullable Class<?> entityClass, String collectionName) {
public Mono<Long> 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));
}
}

View File

@@ -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();
}
}

View File

@@ -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();
}
}

View File

@@ -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<Person> savePersonErrors(Person person) {
TransactionalOperator transactionalOperator = TransactionalOperator.create(manager,
new DefaultTransactionDefinition());
return operations.save(person) //
.<Person> flatMap(it -> Mono.error(new RuntimeException("poof!"))) //
.as(transactionalOperator::transactional);
}
public Mono<Person> savePerson(Person person) {
TransactionalOperator transactionalOperator = TransactionalOperator.create(manager,
new DefaultTransactionDefinition());
return operations.save(person) //
.flatMap(Mono::just) //
.as(transactionalOperator::transactional);
}
public Mono<Long> 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<EventLog> 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<Void> 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"))) //
.<Void> flatMap(it -> Mono.error(new RuntimeException("poof!"))) //
.as(transactionalOperator::transactional);
}
@Transactional
public Mono<Person> declarativeSavePerson(Person person) {
TransactionalOperator transactionalOperator = TransactionalOperator.create(manager,
new DefaultTransactionDefinition());
return operations.save(person) //
.flatMap(Mono::just) //
.as(transactionalOperator::transactional);
}
@Transactional
public Mono<Person> declarativeSavePersonErrors(Person person) {
TransactionalOperator transactionalOperator = TransactionalOperator.create(manager,
new DefaultTransactionDefinition());
return operations.save(person) //
.<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;
}
}

View File

@@ -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<Person> savePerson(Person person) {
TransactionalOperator transactionalOperator = TransactionalOperator.create(manager,
new DefaultTransactionDefinition());
return operations.save(person).<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;
}
}

View File

@@ -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.
*

View File

@@ -17,6 +17,7 @@
* <<mongo.jsonSchema.generated, JSON Schema generation>> from domain types.
* SpEL support in for expressions in `@Indexed`.
* Annotation-based Collation support through `@Document` and `@Query`.
* Declarative reactive transactions using <<mongo.transactions.reactive-tx-manager, @Transactional>>.
[[new-features.2-1-0]]
== What's New in Spring Data MongoDB 2.1

View File

@@ -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<DeleteResult> 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<Void> 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<UpdateResult> 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", …));
};
});