From e585bb949f019fb40d6bec492811bdafe0003a9d Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 19 Dec 2017 15:01:11 -0500 Subject: [PATCH] INT-4368: Twitter - move init to start() JIRA: https://jira.spring.io/browse/INT-4368 Test fails if redis not available. Also fix synchronization of `lastEnqueuedId`. --- .../inbound/AbstractTwitterMessageSource.java | 55 ++++++++++++++----- 1 file changed, 41 insertions(+), 14 deletions(-) diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java index b970626c7f..d43933c101 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java @@ -24,6 +24,7 @@ import java.util.Queue; import java.util.concurrent.LinkedBlockingQueue; import org.springframework.beans.factory.BeanFactory; +import org.springframework.context.Lifecycle; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.context.IntegrationObjectSupport; import org.springframework.integration.core.MessageSource; @@ -57,7 +58,8 @@ import org.springframework.util.StringUtils; * @since 2.0 */ @SuppressWarnings("rawtypes") -abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport implements MessageSource { +abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport implements MessageSource, + Lifecycle { private static final int DEFAULT_PAGE_SIZE = 20; @@ -69,10 +71,10 @@ abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport private final String metadataKey; - private volatile MetadataStore metadataStore; - private final Queue tweets = new LinkedBlockingQueue(); + private volatile MetadataStore metadataStore; + private volatile int prefetchThreshold = 0; private volatile long lastEnqueuedId = -1; @@ -81,6 +83,7 @@ abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport private volatile int pageSize = DEFAULT_PAGE_SIZE; + private volatile boolean running; protected AbstractTwitterMessageSource(Twitter twitter, String metadataKey) { Assert.notNull(twitter, "twitter must not be null"); @@ -132,13 +135,31 @@ abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport } } - String lastId = this.metadataStore.get(this.metadataKey); - // initialize the last status ID from the metadataStore - if (StringUtils.hasText(lastId)) { - this.lastProcessedId = Long.parseLong(lastId); - this.lastEnqueuedId = this.lastProcessedId; - } + } + @Override + public synchronized void start() { + if (!this.running) { + String lastId = this.metadataStore.get(this.metadataKey); + // initialize the last status ID from the metadataStore + if (StringUtils.hasText(lastId)) { + this.lastProcessedId = Long.parseLong(lastId); + synchronized (this.lastEnqueuedIdMonitor) { + this.lastEnqueuedId = this.lastProcessedId; + } + } + this.running = true; + } + } + + @Override + public synchronized void stop() { + this.running = false; + } + + @Override + public synchronized boolean isRunning() { + return this.running; } @Override @@ -169,7 +190,9 @@ abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport long id = this.getIdForTweet(tweet); if (id > this.lastEnqueuedId) { this.tweets.add(tweet); - this.lastEnqueuedId = id; + synchronized (this.lastEnqueuedIdMonitor) { + this.lastEnqueuedId = id; + } } } } @@ -177,9 +200,11 @@ abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport private void refreshTweetQueueIfNecessary() { try { if (this.tweets.size() <= this.prefetchThreshold) { - List tweets = pollForTweets(this.lastEnqueuedId); - if (!CollectionUtils.isEmpty(tweets)) { - enqueueAll(tweets); + synchronized (this.lastEnqueuedIdMonitor) { + List tweets = pollForTweets(this.lastEnqueuedId); + if (!CollectionUtils.isEmpty(tweets)) { + enqueueAll(tweets); + } } } } @@ -222,7 +247,9 @@ abstract class AbstractTwitterMessageSource extends IntegrationObjectSupport synchronized (this) { this.metadataStore.remove(this.metadataKey); this.lastProcessedId = -1L; - this.lastEnqueuedId = -1L; + synchronized (this.lastEnqueuedIdMonitor) { + this.lastEnqueuedId = -1L; + } } }