diff --git a/build.gradle b/build.gradle index 66e1b39f9e..45c8df665d 100644 --- a/build.gradle +++ b/build.gradle @@ -535,6 +535,8 @@ project('spring-integration-twitter') { } compile("javax.activation:activation:$javaxActivationVersion", optional) testCompile project(":spring-integration-test") + testCompile project(":spring-integration-redis") + testCompile project(":spring-integration-redis").sourceSets.test.output } } diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/metadata/RedisMetadataStore.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/metadata/RedisMetadataStore.java new file mode 100644 index 0000000000..b837625b41 --- /dev/null +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/metadata/RedisMetadataStore.java @@ -0,0 +1,68 @@ +/* + * Copyright 2013 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.redis.store.metadata; + +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.data.redis.core.BoundValueOperations; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.integration.store.metadata.MetadataStore; +import org.springframework.util.Assert; + +/** + * Redis implementation of {@link MetadataStore}. Use this {@link MetadataStore} + * to achieve meta-data persistence across application restarts. + * + * @author Gunnar Hillert + * @since 3.0 + */ +public class RedisMetadataStore implements MetadataStore { + + private final StringRedisTemplate redisTemplate; + + /** + * Initializes the {@link RedisTemplate}. + * A {@link StringRedisTemplate} is used with default properties. + * + * @param connectionFactory Must not be null + */ + public RedisMetadataStore(RedisConnectionFactory connectionFactory) { + Assert.notNull(connectionFactory, "'connectionFactory' must not be null."); + this.redisTemplate = new StringRedisTemplate(connectionFactory); + } + + /** + * Persists the provided key and value to Redis. + * + * @param key Must not be null + * @param value Must not be null + */ + public void put(String key, String value) { + Assert.notNull(key, "'key' must not be null."); + Assert.notNull(value, "'value' must not be null."); + BoundValueOperations ops = this.redisTemplate.boundValueOps(key); + ops.set(value); + } + + /** + * Retrieve the persisted value for the provided key. + * + * @param key Must not be null + */ + public String get(String key) { + Assert.notNull(key, "'key' must not be null."); + BoundValueOperations ops = this.redisTemplate.boundValueOps(key); + return ops.get(); + } +} diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/metadata/package-info.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/metadata/package-info.java new file mode 100644 index 0000000000..8e0a51cc48 --- /dev/null +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/metadata/package-info.java @@ -0,0 +1,5 @@ +/** + * Provides support for Redis-based + * {@link org.springframework.integration.store.metadata.MetadataStore}s. + */ +package org.springframework.integration.redis.store.metadata; diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/metadata/RedisMetadataStoreTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/metadata/RedisMetadataStoreTests.java new file mode 100644 index 0000000000..c19191bdb1 --- /dev/null +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/metadata/RedisMetadataStoreTests.java @@ -0,0 +1,145 @@ +/* + * Copyright 2013 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.redis.store.metadata; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.fail; + +import org.junit.Test; +import org.springframework.data.redis.connection.jedis.JedisConnectionFactory; +import org.springframework.data.redis.core.BoundValueOperations; +import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.integration.redis.rules.RedisAvailable; +import org.springframework.integration.redis.rules.RedisAvailableTests; + +/** + * @author Gunnar Hillert + * @since 3.0 + * + */ +public class RedisMetadataStoreTests extends RedisAvailableTests { + + @Test + @RedisAvailable + public void testGetNonExistingKeyValue(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + String retrievedValue = metadataStore.get("does-not-exist"); + assertNull(retrievedValue); + } + + @Test + @RedisAvailable + public void testPersistKeyValue(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + metadataStore.put("RedisMetadataStoreTests-Spring", "Integration"); + + StringRedisTemplate redisTemplate = new StringRedisTemplate(jcf); + BoundValueOperations ops = redisTemplate.boundValueOps("RedisMetadataStoreTests-Spring"); + + assertEquals("Integration", ops.get()); + } + + @Test + @RedisAvailable + public void testGetValueFromMetadataStore(){ + + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + metadataStore.put("RedisMetadataStoreTests-GetValue", "Hello Redis"); + + String retrievedValue = metadataStore.get("RedisMetadataStoreTests-GetValue"); + assertEquals("Hello Redis", retrievedValue); + } + + @Test + @RedisAvailable + public void testPersistEmptyStringToMetadataStore(){ + + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + metadataStore.put("RedisMetadataStoreTests-PersistEmpty", ""); + + String retrievedValue = metadataStore.get("RedisMetadataStoreTests-PersistEmpty"); + assertEquals("", retrievedValue); + } + + @Test + @RedisAvailable + public void testPersistNullStringToMetadataStore(){ + + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + + try { + metadataStore.put("RedisMetadataStoreTests-PersistEmpty", null); + } + catch (IllegalArgumentException e) { + assertEquals("'value' must not be null.", e.getMessage()); + return; + } + + fail("Expected an IllegalArgumentException to be thrown."); + + } + + @Test + @RedisAvailable + public void testPersistWithEmptyKeyToMetadataStore(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + metadataStore.put("", "PersistWithEmptyKey"); + + String retrievedValue = metadataStore.get(""); + assertEquals("PersistWithEmptyKey", retrievedValue); + } + + @Test + @RedisAvailable + public void testPersistWithNullKeyToMetadataStore(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + + try { + metadataStore.put(null, "something"); + } + catch (IllegalArgumentException e) { + assertEquals("'key' must not be null.", e.getMessage()); + return; + } + + fail("Expected an IllegalArgumentException to be thrown."); + } + + @Test + @RedisAvailable + public void testGetValueWithNullKeyFromMetadataStore(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMetadataStore metadataStore = new RedisMetadataStore(jcf); + + try { + metadataStore.get(null); + } + catch (IllegalArgumentException e) { + assertEquals("'key' must not be null.", e.getMessage()); + return; + } + + fail("Expected an IllegalArgumentException to be thrown."); + } +} diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests-context.xml b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests-context.xml new file mode 100644 index 0000000000..ff71c8a29a --- /dev/null +++ b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests-context.xml @@ -0,0 +1,33 @@ + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests.java b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests.java new file mode 100644 index 0000000000..4985dfea89 --- /dev/null +++ b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests.java @@ -0,0 +1,158 @@ +/* + * Copyright 2013 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.twitter.inbound; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; +import static org.mockito.Matchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.util.ArrayList; +import java.util.GregorianCalendar; +import java.util.List; + +import org.junit.Before; +import org.junit.Test; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.ImportResource; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.integration.Message; +import org.springframework.integration.endpoint.SourcePollingChannelAdapter; +import org.springframework.integration.redis.rules.RedisAvailable; +import org.springframework.integration.redis.rules.RedisAvailableTests; +import org.springframework.integration.redis.store.metadata.RedisMetadataStore; +import org.springframework.integration.store.metadata.MetadataStore; +import org.springframework.integration.test.util.TestUtils; +import org.springframework.social.twitter.api.SearchMetadata; +import org.springframework.social.twitter.api.SearchOperations; +import org.springframework.social.twitter.api.SearchResults; +import org.springframework.social.twitter.api.Tweet; +import org.springframework.social.twitter.api.impl.SearchParameters; +import org.springframework.social.twitter.api.impl.TwitterTemplate; + +/** + * @author Gunnar Hillert + * @since 3.0 + */ +public class SearchReceivingMessageSourceWithRedisTests extends RedisAvailableTests { + + private SourcePollingChannelAdapter twitterSearchAdapter; + private RedisConnectionFactory redisConnectionFactory; + private StringRedisTemplate redisTemplate; + + private AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(); + + @Before + public void setup() { + context.register(SearchReceivingMessageSourceWithRedisTestsConfig.class); + context.registerShutdownHook(); + context.refresh(); + + this.redisConnectionFactory = context.getBean(RedisConnectionFactory.class); + this.twitterSearchAdapter = context.getBean(SourcePollingChannelAdapter.class); + this.redisTemplate = new StringRedisTemplate(redisConnectionFactory); + } + + /** + * Verify that a polling operation returns in fact 3 results. + * @throws Exception + */ + @Test + @RedisAvailable + public void testPollForTweetsThreeResultsWithRedisMetadataStore() throws Exception { + + final MetadataStore metadataStore = TestUtils.getPropertyValue(twitterSearchAdapter, "source.metadataStore", MetadataStore.class); + assertTrue("Exptected metadataStore to be an instance of RedisMetadataStore", metadataStore instanceof RedisMetadataStore); + + /* + * The metadataKey is automatically generated. To ensure that we use the + * the correct key, we retrieve it from the adapter. + */ + final String metadataKey = TestUtils.getPropertyValue(twitterSearchAdapter, "source.metadataKey", String.class); + + /* + * As we had to retrieve the metadataKey from the adapter. The metdataStore + * was already invoked and the id retrieved from Redis before we had a chance + * to reset possibly pre-existing values. + * + * Rather than deleting the value, we have to set a value, because "null" values + * returned from the MetadataStore are ignored by the onInit() method in + * the AbstractTwitterMessageSource. */ + redisTemplate.opsForValue().set(metadataKey, "-1"); + assertEquals("-1", redisTemplate.opsForValue().get(metadataKey)); + + final SearchReceivingMessageSource source = TestUtils.getPropertyValue(twitterSearchAdapter, "source", SearchReceivingMessageSource.class); + + /* We need to call onInit() in order to update the id from the metadataStore. */ + source.onInit(); + + final Message message1 = source.receive(); + final Message message2 = source.receive(); + final Message message3 = source.receive(); + + /* We received 3 messages so far. When invoking receive() again the search + * will return again the 3 test Tweets but as we already processed them + * no message (null) is returned. */ + final Message message4 = source.receive(); + + assertNotNull(message1); + assertNotNull(message2); + assertNotNull(message3); + assertNull(message4); + + final String persistedMetadataStoreValue = redisTemplate.opsForValue().get(metadataKey); + assertNotNull(persistedMetadataStoreValue); + assertEquals("3", redisTemplate.opsForValue().get(metadataKey)); + + redisTemplate.delete(metadataKey); + } + + @Configuration + @ImportResource("classpath:org/springframework/integration/twitter/inbound/SearchReceivingMessageSourceWithRedisTests-context.xml") + static class SearchReceivingMessageSourceWithRedisTestsConfig { + + @Bean(name="twitterTemplate") + public TwitterTemplate twitterTemplate() { + final TwitterTemplate twitterTemplate = mock(TwitterTemplate.class); + + final SearchOperations so = mock(SearchOperations.class); + + final Tweet tweet3 = new Tweet(3L, "first", new GregorianCalendar(2013, 2, 20).getTime(), "fromUser", "profileImageUrl", 888L, 999L, "languageCode", "source"); + final Tweet tweet1 = new Tweet(1L, "first", new GregorianCalendar(2013, 0, 20).getTime(), "fromUser", "profileImageUrl", 888L, 999L, "languageCode", "source"); + final Tweet tweet2 = new Tweet(2L, "first", new GregorianCalendar(2013, 1, 20).getTime(), "fromUser", "profileImageUrl", 888L, 999L, "languageCode", "source"); + + final List tweets = new ArrayList(); + + tweets.add(tweet3); + tweets.add(tweet1); + tweets.add(tweet2); + + final SearchResults results = new SearchResults(tweets, new SearchMetadata(111, 111)); + + when(twitterTemplate.searchOperations()).thenReturn(so); + when(twitterTemplate.searchOperations().search(any(SearchParameters.class))).thenReturn(results); + + return twitterTemplate; + } + } +} diff --git a/src/reference/docbook/feed.xml b/src/reference/docbook/feed.xml index 0cb4054364..5215161a09 100644 --- a/src/reference/docbook/feed.xml +++ b/src/reference/docbook/feed.xml @@ -3,33 +3,33 @@ xmlns:xlink="http://www.w3.org/1999/xlink"> Feed Adapter - Spring Integration provides support for Syndication via Feed Adapters + Spring Integration provides support for Syndication via Feed Adapters
Introduction - Web syndication is a form of publishing material such as news stories, press releases, blog posts, and + Web syndication is a form of publishing material such as news stories, press releases, blog posts, and other items typically available on a website but also made available in a feed format such as RSS or ATOM. - Spring integration provides support for Web Syndication via its 'feed' adapter and provides convenient - namespace-based configuration for it. + Spring integration provides support for Web Syndication via its 'feed' adapter and provides convenient + namespace-based configuration for it. To configure the 'feed' namespace, include the following elements within the headers of your XML configuration file: - +
-
+
Feed Inbound Channel Adapter The only adapter that is really needed to provide support for retrieving feeds is an inbound channel adapter. This allows you to subscribe to a particular URL. Below is an example configuration: - - ]]> @@ -38,46 +38,74 @@ xsi:schemaLocation="http://www.springframework.org/schema/integration/feed As news items are retrieved they will be converted to Messages and sent to a channel identified by the channel attribute. - The payload of each message will be a com.sun.syndication.feed.synd.SyndEntry instance. That encapsulates + The payload of each message will be a com.sun.syndication.feed.synd.SyndEntry instance. That encapsulates various data about a news item (content, dates, authors, etc.). - You can also see that the Inbound Feed Channel Adapter is a Polling Consumer. That means you have to + You can also see that the Inbound Feed Channel Adapter is a Polling Consumer. That means you have to provide a poller configuration. However, one important thing you must understand with regard to Feeds is that its inner-workings - are slightly different then most other poling consumers. When an Inbound Feed adapter is started, it does the first poll and - receives a com.sun.syndication.feed.synd.SyndEntryFeed instance. That is an object that contains multiple - SyndEntry objects. Each entry is stored in the local entry queue and is released based on - the value in the max-messages-per-poll attribute such that each Message will contain a single entry. - If during retrieval of the entries from the entry queue the queue had become empty, the adapter will attempt to update + are slightly different then most other poling consumers. When an Inbound Feed adapter is started, it does the first poll and + receives a com.sun.syndication.feed.synd.SyndEntryFeed instance. That is an object that contains multiple + SyndEntry objects. Each entry is stored in the local entry queue and is released based on + the value in the max-messages-per-poll attribute such that each Message will contain a single entry. + If during retrieval of the entries from the entry queue the queue had become empty, the adapter will attempt to update the Feed thereby populating the queue with more entries (SyndEntry instances) if available. Otherwise the next attempt to poll for a feed will be determined by the trigger of the poller (e.g., every 10 seconds in the above configuration). - + Duplicate Entries Polling for a Feed might result in entries that have already been processed - ("I already read that news item, why are you showing it to me again?"). + ("I already read that news item, why are you showing it to me again?"). Spring Integration provides a convenient mechanism to eliminate the need to worry about duplicate entries. - Each feed entry will have a published date field. Every time a new Message is generated and sent, + Each feed entry will have a published date field. Every time a new Message is generated and sent, Spring Integration will store the value of the latest published date in an instance of the org.springframework.integration.store.MetadataStore strategy. The MetadataStore interface is designed to store various types of generic meta-data (e.g., published date of the last feed entry that has been processed) - to help components such as this Feed adapter deal with duplicates. + to help components such as this Feed adapter deal with duplicates. - - The default rule for locating this metadata store is as follows: Spring Integration will look for a bean of type - org.springframework.integration.store.MetadataStore in the ApplicationContext. 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, you may use the - PropertiesPersistingMetadataStore which is backed by a properties file and a properties-persister. - Alternatively, you could provide your own implementation of the MetadataStore interface - (e.g. JdbcMetadataStore) and configure it as bean in the Application Context. - - + The default rule for locating this metadata store is as follows: + Spring Integration will look for a bean of type + org.springframework.integration.store.MetadataStore in + the ApplicationContext. 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 available: + + + PropertiesPersistingMetadataStore + RedisMetadataStore + + + The PropertiesPersistingMetadataStore is backed by + a properties file and a + PropertiesPersister. + + ]]> - -
+ + As of Spring Integration 3.0 a Redis-based + MetadataStore is also available. For + more information regarding the RedisMetadataStore + see . + + + Be careful when using the same Redis instancce across multiple application + contexts as separate Feed adapters may accidentally use the same persisted + key. + + + Alternatively, you could provide your own implementation of the + MetadataStore interface (e.g. JdbcMetadataStore) + and configure it as bean in the Application Context. + +
diff --git a/src/reference/docbook/redis.xml b/src/reference/docbook/redis.xml index 129d716e8a..42dcefeadd 100644 --- a/src/reference/docbook/redis.xml +++ b/src/reference/docbook/redis.xml @@ -223,7 +223,36 @@ rt.setConnectionFactory(redisConnectionFactory);]]> the valueSerializer property of the RedisMessageStore. - +
+ Redis Metadata Store + + As of Spring Integration 3.0 a new Redis-based + MetadataStore + implementation is available. The RedisMetadataStore can + be used to maintain state of a MetadataStore + across application restarts. This new MetadataStore + implementation can be used with adapters such as: + + + Twitter Inbound Adapters + Feed Inbound Channel Adapter + + + In order to instruct these adapters to use the new RedisMetadataStore + 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 RedisMetadataStore. + + + +]]> + + Be careful when using the same Redis instancce across multiple application + contexts as separate adapters may accidentally use the same persisted + key. + +
RedisStore Inbound Channel Adapter diff --git a/src/reference/docbook/twitter.xml b/src/reference/docbook/twitter.xml index 9ef950b015..267f665777 100644 --- a/src/reference/docbook/twitter.xml +++ b/src/reference/docbook/twitter.xml @@ -151,11 +151,22 @@ twitter.oauth.accessTokenSecret=AbRxUAvyNCtqQtxFK8w5ZMtMj20KFhB6o]]>PropertiesPersistingMetadataStore (which is backed by a properties file, and a persister strategy), or you may create your own custom implementation of the MetadataStore interface (e.g., JdbcMetadatStore) and configure it as a bean named 'metadataStore' within the Application Context. + + + As of Spring Integration 3.0 a Redis-based + MetadataStore is available. The + RedisMetadataStore allows you to maintain persisted + metadata across Application Context restarts. For more information see . + + + Be careful when using the same Redis instance across multiple application + contexts as separate Twitter adapters may accidentally use the same persisted + key. + ]]> The Poller that is configured as part of any Inbound Twitter Adapter (see below) will simply poll from this MetadataStore to determine the latest tweet received. -
Inbound Message Channel Adapter diff --git a/src/reference/docbook/whats-new.xml b/src/reference/docbook/whats-new.xml index 074d64a9ad..e6d68b4b5e 100644 --- a/src/reference/docbook/whats-new.xml +++ b/src/reference/docbook/whats-new.xml @@ -132,6 +132,24 @@ For more information see .
+
+ Redis Metadata Store + + A new Redis-based + MetadataStore + implementation was added. The RedisMetadataStore can + be used to maintain state of a MetadataStore + across application restarts. This new MetadataStore + implementation can be used with adapters such as: + + + Twitter Inbound Adapters + Feed Inbound Channel Adapter + + + For more information see . + +