diff --git a/consumer/mongodb-consumer/pom.xml b/consumer/mongodb-consumer/pom.xml
index 2799602e..0ab63380 100644
--- a/consumer/mongodb-consumer/pom.xml
+++ b/consumer/mongodb-consumer/pom.xml
@@ -48,6 +48,11 @@
de.flapdoodle.embed.mongo
test
+
+ org.awaitility
+ awaitility
+ test
+
org.springframework.boot
spring-boot-configuration-processor
diff --git a/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java b/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java
index 421ec335..b667203b 100644
--- a/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java
+++ b/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java
@@ -35,6 +35,7 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
/**
* @author David Turanski
@@ -70,7 +71,11 @@ class MongoDbConsumerApplicationTests {
messages.map(message -> {
mongodbConsumer.accept(message);
return message;
- }).blockLast(Duration.ofSeconds(30));
+
+ }).subscribe();
+
+ await().timeout(Duration.ofSeconds(10))
+ .until(() -> mongoTemplate.findAll(Document.class, properties.getCollection()).count().block() == 3L);
StepVerifier.create(this.mongoTemplate.findAll(Document.class, properties.getCollection())
.sort(Comparator.comparing(d -> d.get("_id").toString())))