diff --git a/build.gradle b/build.gradle index 21bc5fa4cb..ed263f4008 100644 --- a/build.gradle +++ b/build.gradle @@ -458,6 +458,7 @@ project('spring-integration-mongodb') { exclude group: 'org.springframework', module: 'spring-expression' exclude group: 'org.springframework', module: 'spring-tx' } + testCompile "org.slf4j:slf4j-log4j12:$slf4jVersion" } } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java new file mode 100644 index 0000000000..26517bd43c --- /dev/null +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java @@ -0,0 +1,236 @@ +/* + * Copyright 2015 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.integration.mongodb.metadata; + +import java.util.HashMap; +import java.util.Map; + +import org.springframework.dao.DataAccessException; +import org.springframework.data.mongodb.MongoDbFactory; +import org.springframework.data.mongodb.core.CollectionCallback; +import org.springframework.data.mongodb.core.MongoTemplate; +import org.springframework.data.mongodb.core.ScriptOperations; +import org.springframework.data.mongodb.core.query.Criteria; +import org.springframework.data.mongodb.core.query.Query; +import org.springframework.data.mongodb.core.query.Update; +import org.springframework.data.mongodb.core.script.NamedMongoScript; +import org.springframework.integration.metadata.ConcurrentMetadataStore; +import org.springframework.util.Assert; + +import com.mongodb.BasicDBObject; +import com.mongodb.DBCollection; +import com.mongodb.MongoException; + +/** + * MongoDbMetadataStore implementation of {@link ConcurrentMetadataStore}. + * Use this {@link org.springframework.integration.metadata.MetadataStore} to + * achieve meta-data persistence shared across application instances and + * restarts. + * + * @author Senthil Arumugam, Samiraj Panneer Selvam + * @author Artem Bilan + * @since 4.2 + * + */ +public class MongoDbMetadataStore implements ConcurrentMetadataStore { + + private static final String DEFAULT_COLLECTION_NAME = "metadataStore"; + + private static final String ID_FIELD = "_id"; + + private static final String VALUE = "value"; + + private static final String PUT_IF_ABSENT_FUNCTION = + "function putIfAbsent(collection, key, value){ " + + " var alreadyPresent = db[collection].findOne({\"_id\": key}, {\"_id\": 0}); " + + " if(alreadyPresent == null){" + + " db[collection].insert({\"_id\": key, \"value\": value}); " + + " return null; " + + " }" + + " return alreadyPresent;" + + "}"; + + private static final String PUT_IF_ABSENT_SCRIPT_NAME = "metadataStorePutIfAbsent"; + + private final MongoTemplate template; + + private final String collectionName; + + private final ScriptOperations scriptOperations; + + private volatile boolean scriptInitialized; + + /** + * Configure the MongoDbMetadataStore by provided {@link MongoDbFactory} and + * default collection name - {@link #DEFAULT_COLLECTION_NAME}. + * @param factory the mongodb factory + */ + public MongoDbMetadataStore(MongoDbFactory factory) { + this(factory, DEFAULT_COLLECTION_NAME); + } + + /** + * Configure the MongoDbMetadataStore by provided {@link MongoDbFactory} and + * collection name + * @param factory the mongodb factory + * @param collectionName the collection name where it persists the data + */ + public MongoDbMetadataStore(MongoDbFactory factory, String collectionName) { + this(new MongoTemplate(factory), collectionName); + } + + /** + * Configure the MongoDbMetadataStore by provided {@link MongoTemplate} and + * default collection name - {@link #DEFAULT_COLLECTION_NAME}. + * @param template the mongodb template + */ + public MongoDbMetadataStore(MongoTemplate template) { + this(template, DEFAULT_COLLECTION_NAME); + } + + /** + * Configure the MongoDbMetadataStore by provided {@link MongoTemplate} and collection name. + * @param template the mongodb template + * @param collectionName the collection name where it persists the data + */ + public MongoDbMetadataStore(MongoTemplate template, String collectionName) { + Assert.notNull(template, "'template' must not be null."); + Assert.hasText(collectionName, "'collectionName' must not be empty."); + this.template = template; + this.collectionName = collectionName; + this.scriptOperations = template.scriptOps(); + } + + /** + * Store a metadata {@code value} under provided {@code key} to the configured + * {@link #collectionName}. + *

+ * If a document does not exist with the specified {@code key}, the method performs an {@code insert}. + * If a document exists with the specified {@code key}, the method performs an {@code update}. + * @param key the metadata entry key + * @param value the metadata entry value + * @see MongoTemplate#execute(String, CollectionCallback) + * @see DBCollection#save + */ + @Override + public void put(String key, String value) { + Assert.hasText(key, "'key' must not be empty."); + Assert.hasText(value, "'value' must not be empty."); + final Map entry = new HashMap(); + entry.put(ID_FIELD, key); + entry.put(VALUE, value); + this.template.execute(this.collectionName, (CollectionCallback) new CollectionCallback() { + + @Override + public Object doInCollection(DBCollection collection) throws MongoException, DataAccessException { + return collection.save(new BasicDBObject(entry)); + } + + }); + } + + /** + * Get the {@code value} for the provided {@code key} performing {@code findOne} MongoDB operation. + * @param key the metadata entry key + * @return the metadata entry value or null if doesn't exist. + * @see MongoTemplate#findOne(Query, Class, String) + */ + @Override + public String get(String key) { + Assert.hasText(key, "'key' must not be empty."); + Query query = new Query(Criteria.where(ID_FIELD).is(key)); + query.fields().exclude(ID_FIELD); + @SuppressWarnings("unchecked") + Map result = this.template.findOne(query, Map.class, this.collectionName); + return result == null ? null : result.get(VALUE); + + } + + /** + * Remove the metadata entry for the provided {@code key} and return its {@code value}, if any, + * using {@code findAndRemove} MongoDB operation. + * @param key the metadata entry key + * @return the metadata entry value or null if doesn't exist. + * @see MongoTemplate#findAndRemove(Query, Class, String) + */ + @Override + public String remove(String key) { + Assert.hasText(key, "'key' must not be empty."); + Query query = new Query(Criteria.where(ID_FIELD).is(key)); + query.fields().exclude(ID_FIELD); + @SuppressWarnings("unchecked") + Map result = this.template.findAndRemove(query, Map.class, this.collectionName); + return result == null ? null : result.get(VALUE); + } + + /** + * If the specified key is not already associated with a value, associate it with the given value. + * This is equivalent to + *
 {@code
+	 * if (!map.containsKey(key))
+	 *   return map.put(key, value);
+	 * else
+	 *   return map.get(key);
+	 * }
+ * except that the action is performed atomically. + *

+ * Performs the {@code stored} JavaScript function. + * @param key the metadata entry key + * @param value the metadata entry value to store + * @return null if successful, the old value otherwise. + * @see java.util.concurrent.ConcurrentMap#putIfAbsent(Object, Object) + * @see ScriptOperations#call(String, Object...) + */ + @Override + public String putIfAbsent(String key, String value) { + Assert.hasText(key, "'key' must not be empty."); + Assert.hasText(value, "'value' must not be empty."); + if (!this.scriptInitialized) { + synchronized (this) { + if (!this.scriptInitialized) { + this.scriptOperations.register( + new NamedMongoScript(PUT_IF_ABSENT_SCRIPT_NAME, PUT_IF_ABSENT_FUNCTION)); + this.scriptInitialized = true; + } + } + } + BasicDBObject result = + (BasicDBObject) this.scriptOperations.call(PUT_IF_ABSENT_SCRIPT_NAME, this.collectionName, key, value); + return (result == null) ? null : (String) result.get(VALUE); + } + + /** + * Replace an existing metadata entry {@code value} with a new one. Otherwise does nothing. + * Performs {@code updateFirst} if a document for the provided {@code key} and {@code oldValue} + * exists in the {@link #collectionName}. + * @param key the metadata entry key + * @param oldValue the metadata entry old value to replace + * @param newValue the metadata entry new value to put + * @return {@code true} if replace was successful, {@code false} otherwise. + * @see MongoTemplate#updateFirst(Query, Update, String) + */ + @Override + public boolean replace(String key, String oldValue, String newValue) { + Assert.hasText(key, "'key' must not be empty."); + Assert.hasText(oldValue, "'oldValue' must not be empty."); + Assert.hasText(newValue, "'newValue' must not be empty."); + Query query = new Query(Criteria.where(ID_FIELD).is(key).and(VALUE).is(oldValue)); + return this.template.updateFirst(query, Update.update(VALUE, newValue), this.collectionName) + .isUpdateOfExisting(); + } + +} diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/package-info.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/package-info.java new file mode 100644 index 0000000000..3846827ce9 --- /dev/null +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/package-info.java @@ -0,0 +1,4 @@ +/** + * Contains mongodb metadata store related classes + */ +package org.springframework.integration.mongodb.metadata; diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStoreTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStoreTests.java new file mode 100644 index 0000000000..9295e03bf4 --- /dev/null +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStoreTests.java @@ -0,0 +1,152 @@ +/* + * Copyright 2015 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.integration.mongodb.metadata; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + +import org.junit.Before; +import org.junit.Test; + +import org.springframework.data.mongodb.MongoDbFactory; +import org.springframework.data.mongodb.core.MongoTemplate; +import org.springframework.integration.mongodb.rules.MongoDbAvailable; +import org.springframework.integration.mongodb.rules.MongoDbAvailableTests; + +/** + * @author Senthil Arumugam, Samiraj Panneer Selvam + * @since 4.2 + * + */ +public class MongoDbMetadataStoreTests extends MongoDbAvailableTests { + + private final static String DEFAULT_COLLECTION_NAME = "metadataStore"; + + private final String file1 = "/remotepath/filesTodownload/file-1.txt"; + + private final String file1Id = "12345"; + + private MongoDbMetadataStore store = null; + + @Before + public void configure() throws Exception { + final MongoDbFactory mongoDbFactory = this.prepareMongoFactory(DEFAULT_COLLECTION_NAME); + this.store = new MongoDbMetadataStore(mongoDbFactory); + } + + @MongoDbAvailable + @Test + public void testConfigureCustomCollection() throws Exception { + final String collectionName = "testMetadataStore"; + final MongoDbFactory mongoDbFactory = this.prepareMongoFactory(collectionName); + final MongoTemplate template = new MongoTemplate(mongoDbFactory); + store = new MongoDbMetadataStore(template, collectionName); + testBasics(); + } + + @MongoDbAvailable + @Test + public void testConfigureFactory() throws Exception { + final MongoDbFactory mongoDbFactory = this.prepareMongoFactory(DEFAULT_COLLECTION_NAME); + store = new MongoDbMetadataStore(mongoDbFactory); + testBasics(); + } + + @MongoDbAvailable + @Test + public void testConfigureFactorCustomCollection() throws Exception { + final String collectionName = "testMetadataStore"; + final MongoDbFactory mongoDbFactory = this.prepareMongoFactory(collectionName); + store = new MongoDbMetadataStore(mongoDbFactory, collectionName); + testBasics(); + } + + private void testBasics() { + String fileID = store.get(file1); + assertNull(fileID); + + store.put(file1, file1Id); + + fileID = store.get(file1); + assertNotNull(fileID); + assertEquals(file1Id, fileID); + } + + @Test + @MongoDbAvailable + public void testGetFromStore() { + testBasics(); + } + + @Test + @MongoDbAvailable + public void testPutIfAbsent() throws Exception { + String fileID = store.get(file1); + assertNull("Get First time, Key doesnt exists", fileID); + + fileID = store.putIfAbsent(file1, file1Id); + assertNull("Insert First time, Key insertion successful", fileID); + + fileID = store.putIfAbsent(file1, "56789"); + assertNotNull("Key Already Exists - Insertion Failed, for different value", fileID); + assertEquals("Retrieving the Old Value", file1Id, fileID); + + assertEquals("Retrieving the Old Value", file1Id, store.get(file1)); + + } + + @Test + @MongoDbAvailable + public void testRemove() throws Exception { + String fileID = store.remove(file1); + assertNull(fileID); + + fileID = store.putIfAbsent(file1, file1Id); + assertNull(fileID); + + fileID = store.remove(file1); + assertNotNull(fileID); + assertEquals(file1Id, fileID); + + fileID = store.get(file1); + assertNull(fileID); + } + + @Test + @MongoDbAvailable + public void testReplace() throws Exception { + boolean removedValue = store.replace(file1, file1Id, "4567"); + assertFalse(removedValue); + String fileID = store.get(file1); + assertNull(fileID); + + fileID = store.putIfAbsent(file1, file1Id); + assertNull(fileID); + + removedValue = store.replace(file1, file1Id, "4567"); + assertTrue(removedValue); + + fileID = store.get(file1); + assertNotNull(fileID); + assertEquals("4567", fileID); + + } + +} diff --git a/src/reference/asciidoc/meta-data-store.adoc b/src/reference/asciidoc/meta-data-store.adoc index f16a8f970c..2fe96bcf59 100644 --- a/src/reference/asciidoc/meta-data-store.adoc +++ b/src/reference/asciidoc/meta-data-store.adoc @@ -10,11 +10,13 @@ If a component is not directly provided with a reference to a `MetadataStore`, t If one is found then it will be used, otherwise it will create a new instance of `SimpleMetadataStore` which is an in-memory implementation that will only persist metadata within the lifecycle of the currently running Application Context. This means that upon restart you may end up with duplicate entries. -If you need to persist metadata between Application Context restarts, two persistent `MetadataStores` are provided by the framework: +If you need to persist metadata between Application Context restarts, these persistent `MetadataStores` are provided by +the framework: -* PropertiesPersistingMetadataStore +* `PropertiesPersistingMetadataStore` * <> * <> +* <> diff --git a/src/reference/asciidoc/mongodb.adoc b/src/reference/asciidoc/mongodb.adoc index 750b4fa04d..41d63ede75 100644 --- a/src/reference/asciidoc/mongodb.adoc +++ b/src/reference/asciidoc/mongodb.adoc @@ -140,6 +140,38 @@ To configure that scenario, simply extend one message store bean from the other: ---- +[[mongodb-metadata-store]] +==== MongodDB Metadata Store + +As of _Spring Integration 4.2_, a new MongodDB-based `MetadataStore` (<>) implementation is available. +The `MongoDbMetadataStore` can be used to maintain metadata state across application restarts. +This new `MetadataStore` implementation can be used with adapters such as: + +* <> +* <> +* <> +* <> +* <> + +In order to instruct these adapters to use the new `MongoDbMetadataStore`, simply declare a Spring bean using the +bean name *metadataStore*. The _Twitter Inbound Channel Adapter_ and the _Feed Inbound Channel Adapter_ will both +automatically pick up and use the declared `MongoDbMetadataStore`: + +[source,java] +---- +@Bean +public MetadataStore metadataStore(MongoDbFactory factory) { + return new MongoDbMetadataStore(factory, "integrationMetadataStore"); +} +---- + +The `MongoDbMetadataStore` also implements `ConcurrentMetadataStore`, allowing it to be reliably shared across multiple +application instances where only one instance will be allowed to store or modify a key's value. +All these operations are _atomic_ via MongoDB guarantees. For this purpose the `putIfAbsent` operation is implemented +as a _stored_ JavaScript function. Fom more information see +http://docs.spring.io/spring-data/data-mongo/docs/current/reference/html/#mongo.server-side-scripts[Script Operations] +and `MongoDbMetadataStore` JavaDocs. + [[mongodb-inbound-channel-adapter]] === MongoDB Inbound Channel Adapter diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 7092634158..31acafc5f7 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -17,6 +17,11 @@ However, this has some important implications for (some) user environments. For complete details, see <>. +[[x4.2-mongodb-metadata-store]] +==== MongodDB Metadata Store + +The `MongoDbMetadataStore` is now available. For more information, see <>. + [[x4.2-general]] === General Changes