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 index 24e1d84da..7a014309d 100644 --- 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 @@ -303,25 +303,27 @@ public class ReactiveTransactionIntegrationTests { } @Transactional - public Mono declarativeSavePerson(Person person) { + public Flux declarativeSavePerson(Person person) { TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, new DefaultTransactionDefinition()); - return operations.save(person) // - .flatMap(Mono::just) // - .as(transactionalOperator::transactional); + return transactionalOperator.execute(reactiveTransaction -> { + return operations.save(person); + }); } @Transactional - public Mono declarativeSavePersonErrors(Person person) { + public Flux declarativeSavePersonErrors(Person person) { TransactionalOperator transactionalOperator = TransactionalOperator.create(manager, new DefaultTransactionDefinition()); - return operations.save(person) // - . flatMap(it -> Mono.error(new RuntimeException("poof!"))) // - .as(transactionalOperator::transactional); + return transactionalOperator.execute(reactiveTransaction -> { + + return operations.save(person) // + . flatMap(it -> Mono.error(new RuntimeException("poof!"))); + }); } } diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateTests.java index 0701530b6..c921ba72d 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateTests.java @@ -50,6 +50,7 @@ import org.bson.Document; import org.bson.types.ObjectId; import org.junit.After; import org.junit.Before; +import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; @@ -1358,6 +1359,7 @@ public class ReactiveMongoTemplateTests { } @Test // DATAMONGO-1803 + @Ignore("Heavily relying on timing assumptions. Cannot test message resumption properly. Too much race for too little time in between.") public void changeStreamEventsShouldBeEmittedCorrectly() throws InterruptedException { Assumptions.assumeThat(ReplicaSet.required().runsAsReplicaSet()).isTrue(); @@ -1390,6 +1392,7 @@ public class ReactiveMongoTemplateTests { } @Test // DATAMONGO-1803 + @Ignore("Heavily relying on timing assumptions. Cannot test message resumption properly. Too much race for too little time in between.") public void changeStreamEventsShouldBeConvertedCorrectly() throws InterruptedException { Assumptions.assumeThat(ReplicaSet.required().runsAsReplicaSet()).isTrue(); @@ -1422,6 +1425,7 @@ public class ReactiveMongoTemplateTests { } @Test // DATAMONGO-1803 + @Ignore("Heavily relying on timing assumptions. Cannot test message resumption properly. Too much race for too little time in between.") public void changeStreamEventsShouldBeFilteredCorrectly() throws InterruptedException { Assumptions.assumeThat(ReplicaSet.required().runsAsReplicaSet()).isTrue(); @@ -1498,6 +1502,7 @@ public class ReactiveMongoTemplateTests { } @Test // DATAMONGO-1803 + @Ignore("Heavily relying on timing assumptions. Cannot test message resumption properly. Too much race for too little time in between.") public void changeStreamEventsShouldBeResumedCorrectly() throws InterruptedException { Assumptions.assumeThat(ReplicaSet.required().runsAsReplicaSet()).isTrue();