diff --git a/applications/sink/mongodb-sink/README.adoc b/applications/sink/mongodb-sink/README.adoc index 710fe17a..8116167a 100644 --- a/applications/sink/mongodb-sink/README.adoc +++ b/applications/sink/mongodb-sink/README.adoc @@ -20,6 +20,18 @@ The **$$mongodb$$** $$sink$$ has the following options: //tag::configuration-properties[] $$mongodb.consumer.collection$$:: $$The MongoDB collection to store data.$$ *($$String$$, default: `$$$$`)* $$mongodb.consumer.collection-expression$$:: $$The SpEL expression to evaluate MongoDB collection.$$ *($$Expression$$, default: `$$$$`)* +$$spring.data.mongodb.authentication-database$$:: $$Authentication database name.$$ *($$String$$, default: `$$$$`)* +$$spring.data.mongodb.auto-index-creation$$:: $$Whether to enable auto-index creation.$$ *($$Boolean$$, default: `$$$$`)* +$$spring.data.mongodb.database$$:: $$Database name.$$ *($$String$$, default: `$$$$`)* +$$spring.data.mongodb.field-naming-strategy$$:: $$Fully qualified name of the FieldNamingStrategy to use.$$ *($$Class$$, default: `$$$$`)* +$$spring.data.mongodb.grid-fs-database$$:: $$GridFS database name.$$ *($$String$$, default: `$$$$`)* +$$spring.data.mongodb.host$$:: $$Mongo server host. Cannot be set with URI.$$ *($$String$$, default: `$$$$`)* +$$spring.data.mongodb.password$$:: $$Login password of the mongo server. Cannot be set with URI.$$ *($$Character[]$$, default: `$$$$`)* +$$spring.data.mongodb.port$$:: $$Mongo server port. Cannot be set with URI.$$ *($$Integer$$, default: `$$$$`)* +$$spring.data.mongodb.replica-set-name$$:: $$Required replica set name for the cluster. Cannot be set with URI.$$ *($$String$$, default: `$$$$`)* +$$spring.data.mongodb.uri$$:: $$Mongo database URI. Cannot be set with host, port, credentials and replica set name.$$ *($$String$$, default: `$$mongodb://localhost/test$$`)* +$$spring.data.mongodb.username$$:: $$Login user of the mongo server. Cannot be set with URI.$$ *($$String$$, default: `$$$$`)* +$$spring.data.mongodb.uuid-representation$$:: $$Representation to use when converting a UUID to a BSON binary value.$$ *($$UuidRepresentation$$, default: `$$java-legacy$$`, possible values: `UNSPECIFIED`,`STANDARD`,`C_SHARP_LEGACY`,`JAVA_LEGACY`,`PYTHON_LEGACY`)* //end::configuration-properties[] //end::ref-doc[] diff --git a/applications/sink/mongodb-sink/pom.xml b/applications/sink/mongodb-sink/pom.xml index aabeb91d..27aa8433 100644 --- a/applications/sink/mongodb-sink/pom.xml +++ b/applications/sink/mongodb-sink/pom.xml @@ -16,14 +16,24 @@ + + org.springframework.cloud.fn + mongodb-consumer + org.springframework.boot spring-boot-starter-test test - org.springframework.cloud.fn - mongodb-consumer + io.projectreactor + reactor-test + test + + + de.flapdoodle.embed + de.flapdoodle.embed.mongo + test @@ -43,6 +53,7 @@ ${project.version} org.springframework.cloud.fn.consumer.mongo.MongoDbConsumerConfiguration.class + mongodbConsumer diff --git a/applications/sink/mongodb-sink/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties b/applications/sink/mongodb-sink/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties index 62089217..a9e5d009 100644 --- a/applications/sink/mongodb-sink/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties +++ b/applications/sink/mongodb-sink/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties @@ -1 +1,2 @@ -configuration-properties.classes=org.springframework.cloud.fn.consumer.mongo.MongoDbConsumerProperties +configuration-properties.classes=org.springframework.cloud.fn.consumer.mongo.MongoDbConsumerProperties,\ + org.springframework.boot.autoconfigure.mongo.MongoProperties diff --git a/applications/sink/mongodb-sink/src/main/test/org/springframework/cloud/stream/app/sink/mongodb/MongodbSinkTests.java b/applications/sink/mongodb-sink/src/main/test/org/springframework/cloud/stream/app/sink/mongodb/MongodbSinkTests.java new file mode 100644 index 00000000..9b8540a2 --- /dev/null +++ b/applications/sink/mongodb-sink/src/main/test/org/springframework/cloud/stream/app/sink/mongodb/MongodbSinkTests.java @@ -0,0 +1,56 @@ +/* + * Copyright 2020-2020 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.cloud.stream.app.sink.mongodb; + +import org.bson.Document; +import org.junit.jupiter.api.Test; +import reactor.test.StepVerifier; + +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.cloud.fn.consumer.mongo.MongoDbConsumerConfiguration; +import org.springframework.cloud.stream.binder.test.InputDestination; +import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.context.annotation.Import; +import org.springframework.data.mongodb.core.ReactiveMongoTemplate; +import org.springframework.messaging.support.GenericMessage; + +public class MongodbSinkTests { + @Test + void testInsert() { + TestChannelBinderConfiguration.applicationContextRunner(TestApp.class) + .withPropertyValues( + "spring.data.mongodb.port=0", + "spring.data.mongodb.database=test", + "mongodb.consumer.collection=demo", + "spring.cloud.function.definition=mongodbConsumer") + .run(context -> { + ReactiveMongoTemplate template = context.getBean(ReactiveMongoTemplate.class); + InputDestination inputDestination = context.getBean(InputDestination.class); + inputDestination.send(new GenericMessage("{\"foo\":\"bar\"}".getBytes())); + StepVerifier.create(template.findAll(Document.class, "demo")) + .assertNext(document -> document.containsKey("foo")) + .thenCancel() + .verify(); + }); + } + + @SpringBootApplication + @Import(MongoDbConsumerConfiguration.class) + static class TestApp { + + } +} diff --git a/functions/consumer/mongodb-consumer/README.adoc b/functions/consumer/mongodb-consumer/README.adoc index 93e61ead..cc8e9522 100644 --- a/functions/consumer/mongodb-consumer/README.adoc +++ b/functions/consumer/mongodb-consumer/README.adoc @@ -4,11 +4,13 @@ A consumer that allows you to insert records into MongoDB. ## Beans for injection -You can import `MongoDbConsumerConfiguration` in the application and then inject the following bean. +You can import `MongoDbConsumerConfiguration` in the application and then inject one of the following beans. -`Function, Mono> mongodbConsumer` +`Function, Mono> mongodbConsumerFunction` - Allows you to subscribe. -You can use `mongodbConsumer` as a qualifier when injecting. +`Consumer> mongodbConsumer` - Wraps the function as a Consumer with no-op subscriber. + +You can use `mongodbConsumer` or `mongodbConsumerFunction` as a qualifier when injecting. The return value from the function can be ignored as this is used as a consumer to send records to MongoDB. @@ -18,7 +20,7 @@ All configuration properties are prefixed with `mongodb.consumer`. For more information on the various options available, please see link:src/main/java/org/springframework/cloud/fn/consumer/mongo/MongoDBConsumerProperties.java[MongoDBConsumerProperties]. -## Tests +## Examples See this link:src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDBConsumerApplicationTests.java[test suite] for the various ways, this consumer is used. diff --git a/functions/consumer/mongodb-consumer/src/main/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerConfiguration.java b/functions/consumer/mongodb-consumer/src/main/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerConfiguration.java index e2315858..ff270ace 100644 --- a/functions/consumer/mongodb-consumer/src/main/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerConfiguration.java +++ b/functions/consumer/mongodb-consumer/src/main/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerConfiguration.java @@ -16,6 +16,7 @@ package org.springframework.cloud.fn.consumer.mongo; +import java.util.function.Consumer; import java.util.function.Function; import reactor.core.publisher.Mono; @@ -52,7 +53,12 @@ public class MongoDbConsumerConfiguration { } @Bean - public Function, Mono> mongodbConsumer(ReactiveMessageHandler mongoConsumerMessageHandler) { + public Consumer> mongodbConsumer(Function, Mono> mongodbConsumerFunction) { + return message -> mongodbConsumerFunction.apply(message).subscribe(); + } + + @Bean + public Function, Mono> mongodbConsumerFunction(ReactiveMessageHandler mongoConsumerMessageHandler) { return mongoConsumerMessageHandler::handleMessage; } diff --git a/functions/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java b/functions/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java index 3a761af7..3290517d 100644 --- a/functions/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java +++ b/functions/consumer/mongodb-consumer/src/test/java/org/springframework/cloud/fn/consumer/mongo/MongoDbConsumerApplicationTests.java @@ -20,12 +20,11 @@ import java.time.Duration; import java.util.Comparator; import java.util.HashMap; import java.util.Map; -import java.util.function.Function; +import java.util.function.Consumer; import org.bson.Document; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; import reactor.test.StepVerifier; import org.springframework.beans.factory.annotation.Autowired; @@ -42,14 +41,14 @@ import static org.assertj.core.api.Assertions.assertThat; */ @SpringBootTest(properties = { "spring.data.mongodb.port=0", - "mongodb.consumer.collection=testing"}) + "mongodb.consumer.collection=testing" }) class MongoDbConsumerApplicationTests { @Autowired private MongoDbConsumerProperties properties; @Autowired - private Function, Mono> mongoDbConsumer; + private Consumer> mongodbConsumer; @Autowired private ReactiveMongoTemplate mongoTemplate; @@ -66,10 +65,12 @@ class MongoDbConsumerApplicationTests { Flux> messages = Flux.just( new GenericMessage<>(data1), new GenericMessage<>(data2), - new GenericMessage<>("{\"my_data\": \"THE DATA\"}") - ); + new GenericMessage<>("{\"my_data\": \"THE DATA\"}")); - messages.flatMap(mongoDbConsumer::apply).blockLast(Duration.ofSeconds(10)); + messages.map(message -> { + mongodbConsumer.accept(message); + return message; + }).blockLast(Duration.ofSeconds(10)); StepVerifier.create(this.mongoTemplate.findAll(Document.class, properties.getCollection()) .sort(Comparator.comparing(d -> d.get("_id").toString())))