diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd b/org.springframework.integration/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd index 1340f694ae..7ff34755c4 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd @@ -226,9 +226,9 @@ + - diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java b/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java index cd0e2224e3..4f488c2107 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/AbstractDispatcher.java @@ -23,7 +23,9 @@ import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.core.task.TaskExecutor; +import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageConsumer; +import org.springframework.integration.message.MessageRejectedException; /** * Base class for {@link MessageDispatcher} implementations. @@ -64,4 +66,21 @@ public abstract class AbstractDispatcher implements MessageDispatcher { return this.getClass().getSimpleName() + " with subscribers: " + this.subscribers; } + /** + * Convenience method available for subclasses. Returns 'true' unless a + * "Selective Consumer" throws a {@link MessageRejectedException}. + */ + protected boolean sendMessageToConsumer(Message message, MessageConsumer consumer) { + try { + consumer.onMessage(message); + return true; + } + catch (MessageRejectedException e) { + if (logger.isDebugEnabled()) { + logger.debug("Consumer '" + consumer + "' rejected Message, continuing with other subscribers if available.", e); + } + return false; + } + } + } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java b/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java index 044ea55bd7..c0b211e22c 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/BroadcastingDispatcher.java @@ -56,12 +56,12 @@ public class BroadcastingDispatcher extends AbstractDispatcher { if (executor != null) { executor.execute(new Runnable() { public void run() { - consumer.onMessage(messageToSend); + BroadcastingDispatcher.this.sendMessageToConsumer(messageToSend, consumer); } }); } else { - consumer.onMessage(messageToSend); + this.sendMessageToConsumer(messageToSend, consumer); } } return true; diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/SimpleDispatcher.java b/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/SimpleDispatcher.java index f1608d627d..8c868abaae 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/SimpleDispatcher.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/dispatcher/SimpleDispatcher.java @@ -42,16 +42,10 @@ public class SimpleDispatcher extends AbstractDispatcher { int rejectedExceptionCount = 0; for (MessageConsumer consumer : this.subscribers) { count++; - try { - consumer.onMessage(message); + if (this.sendMessageToConsumer(message, consumer)) { return true; } - catch (MessageRejectedException e) { - rejectedExceptionCount++; - if (logger.isDebugEnabled()) { - logger.debug("Consumer '" + consumer + "' rejected Message, continuing with other subscribers if available.", e); - } - } + rejectedExceptionCount++; } if (rejectedExceptionCount == count) { throw new MessageRejectedException(message, "All of dispatcher's subscribers rejected Message.");