diff --git a/mongodb/change-streams/pom.xml b/mongodb/change-streams/pom.xml
index 1b409e86..df181492 100644
--- a/mongodb/change-streams/pom.xml
+++ b/mongodb/change-streams/pom.xml
@@ -1,65 +1,55 @@
- 4.0.0
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
+ 4.0.0
-
- org.springframework.data.examples
- spring-data-mongodb-examples
- 2.0.0.BUILD-SNAPSHOT
-
+
+ org.springframework.data.examples
+ spring-data-mongodb-examples
+ 2.0.0.BUILD-SNAPSHOT
+
- spring-data-mongodb-change-streams
- Spring Data MongoDB - Change Streams
+ spring-data-mongodb-change-streams
+ Spring Data MongoDB - Change Streams
-
+
+ 3.8.0-beta2
+ 1.9.0-beta1
+ Lovelace-M3
+
-
- org.springframework.boot
- spring-boot-starter-data-mongodb-reactive
-
-
- org.mongodb
- mongodb-driver
-
-
-
+
-
- org.mongodb
- mongo-java-driver
- 3.6.2
-
+
-
- org.mongodb
- mongodb-driver-async
- 3.6.2
-
+
+ spring-data-next
+
+ Lovelace-M3
+
+
-
- org.mongodb
- mongodb-driver-reactivestreams
- 1.7.0
-
+
-
- org.springframework.data
- spring-data-mongodb
- 2.1.0.M1
-
+
-
- io.projectreactor
- reactor-test
- test
-
+
+ org.springframework.boot
+ spring-boot-starter-data-mongodb-reactive
+
-
- de.flapdoodle.embed
- de.flapdoodle.embed.mongo
- 2.0.1
-
+
+ io.projectreactor
+ reactor-test
+ test
+
-
+
+ ${project.groupId}
+ spring-data-mongodb-example-utils
+ ${project.version}
+ test
+
+
+
diff --git a/mongodb/change-streams/src/main/java/example/springdata/mongodb/CollectingMessageListener.java b/mongodb/change-streams/src/main/java/example/springdata/mongodb/CollectingMessageListener.java
index 49cdd6fd..e138f192 100644
--- a/mongodb/change-streams/src/main/java/example/springdata/mongodb/CollectingMessageListener.java
+++ b/mongodb/change-streams/src/main/java/example/springdata/mongodb/CollectingMessageListener.java
@@ -17,6 +17,7 @@ package example.springdata.mongodb;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.atomic.AtomicInteger;
import org.springframework.data.mongodb.core.messaging.Message;
import org.springframework.data.mongodb.core.messaging.MessageListener;
@@ -26,10 +27,12 @@ import org.springframework.data.mongodb.core.messaging.MessageListener;
* and printing the message itself to the console.
*
* @author Christoph Strobl
+ * @author Mark Paluch
*/
class CollectingMessageListener implements MessageListener {
private final BlockingDeque> messages = new LinkedBlockingDeque<>();
+ private final AtomicInteger count = new AtomicInteger();
/*
* (non-Javadoc)
@@ -38,12 +41,22 @@ class CollectingMessageListener implements MessageListener {
@Override
public void onMessage(Message message) {
+
System.out.println(String.format("Received Message in collection %s.\n\trawsource: %s\n\tconverted: %s",
message.getProperties().getCollectionName(), message.getRaw(), message.getBody()));
+
+ count.incrementAndGet();
messages.add(message);
}
int messageCount() {
- return messages.size();
+ return count.get();
+ }
+
+ void awaitNextMessages(int count) throws InterruptedException {
+
+ for (int i = 0; i < count; i++) {
+ messages.take();
+ }
}
}
diff --git a/mongodb/change-streams/src/test/java/example/springdata/mongodb/ChangeStreamsTests.java b/mongodb/change-streams/src/test/java/example/springdata/mongodb/ChangeStreamsTests.java
index 7f3ca223..32b326bb 100644
--- a/mongodb/change-streams/src/test/java/example/springdata/mongodb/ChangeStreamsTests.java
+++ b/mongodb/change-streams/src/test/java/example/springdata/mongodb/ChangeStreamsTests.java
@@ -21,18 +21,19 @@ import static org.springframework.data.mongodb.core.query.Criteria.*;
import static org.springframework.data.mongodb.core.query.Query.*;
import static org.springframework.data.mongodb.core.query.Update.*;
-import reactor.core.Disposable;
+import example.springdata.mongodb.util.EmbeddedMongo;
+import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import java.time.Duration;
-import java.util.concurrent.BlockingQueue;
-import java.util.concurrent.LinkedBlockingQueue;
import org.bson.Document;
+import org.junit.ClassRule;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.autoconfigure.mongo.embedded.EmbeddedMongoAutoConfiguration;
import org.springframework.boot.test.autoconfigure.data.mongo.DataMongoTest;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -40,6 +41,8 @@ import org.springframework.data.mongodb.core.ChangeStreamEvent;
import org.springframework.data.mongodb.core.ChangeStreamOptions;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.ReactiveMongoOperations;
+import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
+import org.springframework.data.mongodb.core.SimpleReactiveMongoDatabaseFactory;
import org.springframework.data.mongodb.core.messaging.ChangeStreamRequest;
import org.springframework.data.mongodb.core.messaging.DefaultMessageListenerContainer;
import org.springframework.data.mongodb.core.messaging.Message;
@@ -48,23 +51,26 @@ import org.springframework.data.mongodb.core.messaging.Subscription;
import org.springframework.test.context.junit4.SpringRunner;
import com.mongodb.client.model.changestream.ChangeStreamDocument;
+import com.mongodb.reactivestreams.client.MongoClients;
/**
* A simple Test demonstrating required {@link Configuration} for consumption of MongoDB
- * Change Streams using the sync and reactivestreams java
+ * Change Streams using the sync and Reactive Streams Java
* driver.
- *
- * We currently need to have both tests in one class due to some shutdown error in flapdoodle which lets us start the
- * process only once.
*
* @author Christoph Strobl
+ * @author Mark Paluch
*/
@RunWith(SpringRunner.class)
@DataMongoTest
-public class ChangeStreamsTests extends ReplicaSetInitiatingTest {
+public class ChangeStreamsTests {
+
+ public static @ClassRule EmbeddedMongo replSet = EmbeddedMongo.replSet().configure();
@Autowired MessageListenerContainer container; // for imperative style
+ @Autowired MongoTemplate template;
+
@Autowired ReactiveMongoOperations reactiveTemplate; // for reactive style
Person gabriel = new Person("Gabriel", "Lorca", 30);
@@ -74,10 +80,30 @@ public class ChangeStreamsTests extends ReplicaSetInitiatingTest {
/**
* Configuration? Yes we need a bit of it - Do not worry, it won't be much!
*/
- @Configuration
- @EnableAutoConfiguration
+ @SpringBootApplication(exclude = EmbeddedMongoAutoConfiguration.class)
static class Config {
+ /**
+ * Configure {@link SimpleMongoDbFactory} pointing to the embedded MongoDB connection.
+ *
+ * @return a new {@link SimpleReactiveMongoDatabaseFactory}.
+ */
+ @Bean
+ SimpleMongoDbFactory mongoDbFactory() {
+ return new SimpleMongoDbFactory(replSet.getMongoClient(), "changestreams");
+ }
+
+ /**
+ * Configure {@link SimpleReactiveMongoDatabaseFactory} pointing to the embedded MongoDB connection.
+ *
+ * @return a new {@link SimpleReactiveMongoDatabaseFactory}.
+ */
+ @Bean
+ SimpleReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory() {
+ return new SimpleReactiveMongoDatabaseFactory(MongoClients.create(replSet.getConnectionString()),
+ "changestreams");
+ }
+
/**
* Since listening to a Change Stream using the sync
* MongoDB Java Driver is a blocking class, we need to move load to another {@link Thread} by simply using a
@@ -105,19 +131,18 @@ public class ChangeStreamsTests extends ReplicaSetInitiatingTest {
/**
* Use the {@link MessageListenerContainer} registered within the
- * {@link org.springframework.context.ApplicationContext} to subscribe to a MongoDB Change Stream. Events published
- * via {@link com.mongodb.client.ChangeStreamIterable} are passed to the
+ * {@link org.springframework.context.ApplicationContext} to subscribe to MongoDB Change Streams. Events published via
+ * {@link com.mongodb.client.ChangeStreamIterable} are passed to the
* {@link org.springframework.data.mongodb.core.messaging.MessageListener#onMessage(Message) MessageListener}.
*/
@Test
public void imperativeChangeEvents() throws InterruptedException {
- CollectingMessageListener, Person> messageListener = new CollectingMessageListener();
+ CollectingMessageListener, Person> messageListener = new CollectingMessageListener<>();
- ChangeStreamRequest request = ChangeStreamRequest.builder() //
+ ChangeStreamRequest request = ChangeStreamRequest.builder(messageListener) //
.collection("person") //
.filter(newAggregation(match(where("operationType").is("insert")))) // we are only interested in inserts
- .publishTo(messageListener) //
.build();
Subscription subscription = container.register(request, Person.class);
@@ -126,7 +151,7 @@ public class ChangeStreamsTests extends ReplicaSetInitiatingTest {
template.save(gabriel);
template.save(ash);
- Thread.sleep(200);
+ messageListener.awaitNextMessages(2);
assertThat(messageListener.messageCount()).isEqualTo(2); // first two insert events, so far so good
@@ -141,46 +166,46 @@ public class ChangeStreamsTests extends ReplicaSetInitiatingTest {
template.save(michael);
- Thread.sleep(200);
+ messageListener.awaitNextMessages(1);
assertThat(messageListener.messageCount()).isEqualTo(3); // there we go, all events received.
}
/**
- * Use the {@link reactor.core.publisher.Flux} to subscribe to a MongoDB Change Stream.
+ * Use a {@link reactor.core.publisher.Flux} to subscribe to MongoDB Change Streams.
*/
@Test
- public void reactiveChangeEvents() throws InterruptedException {
+ public void reactiveChangeEvents() {
- BlockingQueue> documents = new LinkedBlockingQueue<>(100);
+ Flux> changeStream = reactiveTemplate.changeStream(
+ newAggregation(match(where("operationType").is("insert"))), Person.class, ChangeStreamOptions.empty(),
+ "person");
- Disposable disposable = reactiveTemplate.changeStream(newAggregation(match(where("operationType").is("insert"))),
- Person.class, ChangeStreamOptions.empty(), "person").doOnNext(documents::add).subscribe();
+ StepVerifier.create(changeStream) //
+ .expectSubscription() //
+ .expectNoEvent(Duration.ofMillis(200)) // wait till change streams becomes active
- Thread.sleep(200); // wait till the subscription becomes active
+ // Save documents and await their change events
+ .then(() -> {
+ StepVerifier.create(reactiveTemplate.save(gabriel)).expectNextCount(1).verifyComplete();
+ StepVerifier.create(reactiveTemplate.save(ash)).expectNextCount(1).verifyComplete();
+ }).expectNextCount(2) //
- StepVerifier.create(reactiveTemplate.save(gabriel)).expectNextCount(1).verifyComplete();
- StepVerifier.create(reactiveTemplate.save(ash)).expectNextCount(1).verifyComplete();
+ // Update a document
+ .then(() -> {
- Thread.sleep(200);
+ StepVerifier.create(reactiveTemplate.update(Person.class) //
+ .matching(query(where("id").is(ash.getId()))) //
+ .apply(update("age", 40)) //
+ .first()).expectNextCount(1).verifyComplete();
+ }).expectNoEvent(Duration.ofMillis(200)) // updates are skipped
- assertThat(documents.size()).isEqualTo(2); // first two insert events, so far so good
+ // Save another document and await its change event
+ .then(() -> {
+ StepVerifier.create(reactiveTemplate.save(michael)).expectNextCount(1).verifyComplete();
+ }).expectNextCount(1) // there we go, all events received.
- StepVerifier.create(reactiveTemplate.update(Person.class) //
- .matching(query(where("id").is(ash.getId()))) //
- .apply(update("age", 40)) //
- .first()).expectNextCount(1).verifyComplete();
-
- Thread.sleep(200);
-
- assertThat(documents.size()).isEqualTo(2); // updates are skipped
-
- StepVerifier.create(reactiveTemplate.save(michael)).expectNextCount(1).verifyComplete();
-
- Thread.sleep(200); // just give it some time to link receive all events
-
- assertThat(documents.size()).isEqualTo(3); // there we go, all events received.
-
- disposable.dispose();
+ .thenCancel() // change streams are infinite streams, at some point we need to unsubscribe
+ .verify();
}
}
diff --git a/mongodb/change-streams/src/test/java/example/springdata/mongodb/ReplicaSetInitiatingTest.java b/mongodb/change-streams/src/test/java/example/springdata/mongodb/ReplicaSetInitiatingTest.java
deleted file mode 100644
index a58edc38..00000000
--- a/mongodb/change-streams/src/test/java/example/springdata/mongodb/ReplicaSetInitiatingTest.java
+++ /dev/null
@@ -1,70 +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 example.springdata.mongodb;
-
-import java.util.Arrays;
-
-import org.bson.Document;
-import org.hamcrest.number.IsCloseTo;
-import org.junit.Assume;
-import org.junit.Before;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.beans.factory.annotation.Value;
-import org.springframework.data.mongodb.core.MongoTemplate;
-
-import com.mongodb.MongoClient;
-import com.mongodb.MongoClientOptions;
-import com.mongodb.ServerAddress;
-import com.mongodb.client.MongoDatabase;
-
-/**
- * Just a simple base Class initiating the required replica set.
- *
- * @author Christoph Strobl
- */
-abstract class ReplicaSetInitiatingTest {
-
- @Autowired MongoTemplate template;
- @Value("${spring.data.mongodb.port}") Integer port;
-
- /**
- * Before we can actually continue we need to make sure that the Server works as a single node replica set which is
- * required for Change Streams.
- */
- @Before
- public void setUp() {
-
- MongoClient mongo = new MongoClient(new ServerAddress("localhost:" + port),
- MongoClientOptions.builder().connectTimeout(10).build());
-
- MongoDatabase admin = mongo.getDatabase("admin");
- Document status = admin.runCommand(new Document("serverStatus", "1"));
-
- if (!status.get("repl", Document.class).get("ismaster", Boolean.class)) {
-
- Document result = admin.runCommand(new Document("replSetInitiate", new Document("_id", "spring-data-examples")
- .append("members", Arrays.asList(new Document("_id", 0).append("host", "localhost:" + port)))));
-
- Assume.assumeThat(result.getDouble("ok"), IsCloseTo.closeTo(1.0D, 0D));
- }
-
- mongo.close();
-
- template.dropCollection(Person.class);
- template.createCollection(Person.class);
- }
-
-}
diff --git a/mongodb/change-streams/src/test/resources/application.properties b/mongodb/change-streams/src/test/resources/application.properties
deleted file mode 100644
index bbf03ded..00000000
--- a/mongodb/change-streams/src/test/resources/application.properties
+++ /dev/null
@@ -1,4 +0,0 @@
-spring.mongodb.embedded.version=3.6.0
-spring.mongodb.embedded.features=ONLY_WITH_SSL,ONLY_64BIT,NO_HTTP_INTERFACE_ARG
-spring.mongodb.embedded.storage.repl-set-name=spring-data-examples
-spring.data.mongodb.port=57143
\ No newline at end of file
diff --git a/mongodb/change-streams/src/test/resources/logback.xml b/mongodb/change-streams/src/test/resources/logback.xml
new file mode 100644
index 00000000..ef765386
--- /dev/null
+++ b/mongodb/change-streams/src/test/resources/logback.xml
@@ -0,0 +1,17 @@
+
+
+
+
+
+ %d %5p %40.40c:%4L - %m%n
+
+
+
+
+
+
+
+
+
+
+
diff --git a/mongodb/transactions/pom.xml b/mongodb/transactions/pom.xml
index c06f8297..481ebbb0 100644
--- a/mongodb/transactions/pom.xml
+++ b/mongodb/transactions/pom.xml
@@ -37,6 +37,12 @@
spring-boot-starter-data-mongodb-reactive
+
+ io.projectreactor
+ reactor-test
+ test
+
+
${project.groupId}
spring-data-mongodb-example-utils
@@ -44,13 +50,6 @@
test
-
- io.projectreactor
- reactor-test
- 3.1.7.RELEASE
- test
-
-