Added ChannelPollingMessageDispatcher for simplification (no need to set a ChannelPollingMessageRetriever on a dispatcher). The DefaultMessageDispatcher has been renamed BasePollingMessageDispatcher.
This commit is contained in:
@@ -27,14 +27,14 @@ import org.springframework.integration.message.Message;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* The default implementation of {@link MessageDispatcher}. If
|
||||
* The base implementation of a polling {@link MessageDispatcher}. If
|
||||
* {@link #broadcast} is set to <code>false</code> (the default), each message
|
||||
* will be sent to a single {@link MessageHandler}. Otherwise, each
|
||||
* retrieved {@link Message} will be sent to all handlers.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class DefaultMessageDispatcher extends AbstractMessageDispatcher {
|
||||
public class BasePollingMessageDispatcher extends AbstractMessageDispatcher {
|
||||
|
||||
private boolean broadcast = false;
|
||||
|
||||
@@ -45,7 +45,7 @@ public class DefaultMessageDispatcher extends AbstractMessageDispatcher {
|
||||
private boolean shouldFailOnRejectionLimit = true;
|
||||
|
||||
|
||||
public DefaultMessageDispatcher(MessageRetriever retriever) {
|
||||
public BasePollingMessageDispatcher(MessageRetriever retriever) {
|
||||
super(retriever);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,36 @@
|
||||
/*
|
||||
* Copyright 2002-2007 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.bus;
|
||||
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
|
||||
/**
|
||||
* A {@link MessageDispatcher} that polls a {@link MessageChannel}.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class ChannelPollingMessageDispatcher extends BasePollingMessageDispatcher {
|
||||
|
||||
public ChannelPollingMessageDispatcher(MessageChannel channel, int period) {
|
||||
this(channel, ConsumerPolicy.newPollingPolicy(period));
|
||||
}
|
||||
|
||||
public ChannelPollingMessageDispatcher(MessageChannel channel, ConsumerPolicy policy) {
|
||||
super(new ChannelPollingMessageRetriever(channel, policy));
|
||||
}
|
||||
|
||||
}
|
||||
@@ -260,8 +260,7 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif
|
||||
|
||||
private void doActivate(MessageChannel channel, MessageHandler handler, ConsumerPolicy policy) {
|
||||
PooledMessageHandler pooledHandler = new PooledMessageHandler(handler, policy.getConcurrency(), policy.getMaxConcurrency());
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setRejectionLimit(policy.getRejectionLimit());
|
||||
dispatcher.setRetryInterval(policy.getRetryInterval());
|
||||
dispatcher.addHandler(pooledHandler);
|
||||
|
||||
@@ -22,9 +22,8 @@ import java.io.ByteArrayOutputStream;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.bus.ChannelPollingMessageRetriever;
|
||||
import org.springframework.integration.bus.ConsumerPolicy;
|
||||
import org.springframework.integration.bus.DefaultMessageDispatcher;
|
||||
import org.springframework.integration.bus.ChannelPollingMessageDispatcher;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
@@ -42,8 +41,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream);
|
||||
adapter.setChannel(channel);
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
channel.send(new StringMessage("foo"));
|
||||
@@ -60,8 +58,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream);
|
||||
adapter.setChannel(channel);
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
channel.send(new StringMessage("foo"));
|
||||
@@ -82,8 +79,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
adapter.setChannel(channel);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
channel.send(new StringMessage("foo"));
|
||||
@@ -105,8 +101,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
adapter.setChannel(channel);
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
policy.setMaxMessagesPerTask(2);
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
channel.send(new StringMessage("foo"));
|
||||
@@ -126,8 +121,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
policy.setReceiveTimeout(0);
|
||||
policy.setMaxMessagesPerTask(10);
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
channel.send(new StringMessage("foo"));
|
||||
@@ -145,8 +139,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream);
|
||||
adapter.setChannel(channel);
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
TestObject testObject = new TestObject("foo");
|
||||
@@ -166,8 +159,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
policy.setReceiveTimeout(0);
|
||||
policy.setMaxMessagesPerTask(2);
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
@@ -189,8 +181,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy();
|
||||
policy.setReceiveTimeout(0);
|
||||
policy.setMaxMessagesPerTask(2);
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(adapter);
|
||||
dispatcher.start();
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
|
||||
@@ -19,7 +19,6 @@ package org.springframework.integration.bus;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.RejectedExecutionException;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
@@ -35,7 +34,7 @@ import org.springframework.integration.message.selector.PayloadTypeSelector;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class DefaultMessageDispatcherTests {
|
||||
public class ChannelPollingMessageDispatcherTests {
|
||||
|
||||
@Test
|
||||
public void testNonBroadcastingDispatcherSendsToExactlyOneEndpoint() throws InterruptedException {
|
||||
@@ -47,8 +46,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1));
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1));
|
||||
dispatcher.start();
|
||||
@@ -67,8 +65,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setBroadcast(true);
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1));
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1));
|
||||
@@ -90,8 +87,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) {
|
||||
@Override
|
||||
public void start() {
|
||||
@@ -118,8 +114,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setBroadcast(true);
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1));
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) {
|
||||
@@ -140,8 +135,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
assertEquals(0, dispatcher.dispatch());
|
||||
}
|
||||
|
||||
@@ -157,8 +151,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setBroadcast(true);
|
||||
dispatcher.setRejectionLimit(2);
|
||||
dispatcher.setRetryInterval(3);
|
||||
@@ -186,8 +179,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setBroadcast(true);
|
||||
dispatcher.setRejectionLimit(2);
|
||||
dispatcher.setRetryInterval(3);
|
||||
@@ -217,8 +209,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setRejectionLimit(2);
|
||||
dispatcher.setRetryInterval(3);
|
||||
dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) {
|
||||
@@ -249,8 +240,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setRejectionLimit(2);
|
||||
dispatcher.setRetryInterval(3);
|
||||
dispatcher.setShouldFailOnRejectionLimit(false);
|
||||
@@ -290,8 +280,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setRejectionLimit(2);
|
||||
dispatcher.setRetryInterval(3);
|
||||
dispatcher.setShouldFailOnRejectionLimit(false);
|
||||
@@ -342,8 +331,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setBroadcast(true);
|
||||
dispatcher.setRejectionLimit(5);
|
||||
dispatcher.setRetryInterval(3);
|
||||
@@ -387,8 +375,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
PooledMessageHandler executor1 = new PooledMessageHandler(endpoint1, 1, 1);
|
||||
PooledMessageHandler executor2 = new PooledMessageHandler(endpoint2, 1, 1);
|
||||
executor1.addMessageSelector(new PayloadTypeSelector(Integer.class));
|
||||
@@ -415,8 +402,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
PooledMessageHandler executor1 = new PooledMessageHandler(endpoint1, 1, 1) {
|
||||
@Override
|
||||
public Message handle(Message<?> message) {
|
||||
@@ -458,8 +444,7 @@ public class DefaultMessageDispatcherTests {
|
||||
ConsumerPolicy policy = new ConsumerPolicy();
|
||||
SimpleChannel channel = new SimpleChannel();
|
||||
channel.send(new StringMessage(1, "test"));
|
||||
MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy);
|
||||
DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever);
|
||||
ChannelPollingMessageDispatcher dispatcher = new ChannelPollingMessageDispatcher(channel, policy);
|
||||
dispatcher.setBroadcast(true);
|
||||
PooledMessageHandler executor1 = new PooledMessageHandler(endpoint1, 1, 1);
|
||||
PooledMessageHandler executor2 = new PooledMessageHandler(endpoint2, 1, 1);
|
||||
Reference in New Issue
Block a user