SimpleChannel now creates a SynchronousQueue if the capacity is less than or equal to 0. SplitterMessageHandlerAdapter exposes a configurable 'sendTimeout' property.
This commit is contained in:
@@ -20,6 +20,7 @@ import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.BlockingQueue;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.SynchronousQueue;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.springframework.integration.dispatcher.DispatcherPolicy;
|
||||
@@ -46,7 +47,12 @@ public class SimpleChannel extends AbstractMessageChannel {
|
||||
*/
|
||||
public SimpleChannel(int capacity, DispatcherPolicy dispatcherPolicy) {
|
||||
super((dispatcherPolicy != null) ? dispatcherPolicy : new DispatcherPolicy());
|
||||
this.queue = new LinkedBlockingQueue<Message<?>>(capacity);
|
||||
if (capacity > 0) {
|
||||
this.queue = new LinkedBlockingQueue<Message<?>>(capacity);
|
||||
}
|
||||
else {
|
||||
this.queue = new SynchronousQueue<Message<?>>(true);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -47,6 +47,8 @@ public class SplitterMessageHandlerAdapter extends AbstractMessageHandlerAdapter
|
||||
|
||||
private ChannelRegistry channelRegistry;
|
||||
|
||||
private long sendTimeout = -1;
|
||||
|
||||
|
||||
public SplitterMessageHandlerAdapter(Object object, Method method, Map<String, ?> attributes) {
|
||||
Assert.notNull(object, "'object' must not be null");
|
||||
@@ -63,6 +65,10 @@ public class SplitterMessageHandlerAdapter extends AbstractMessageHandlerAdapter
|
||||
this.channelRegistry = channelRegistry;
|
||||
}
|
||||
|
||||
public void setSendTimeout(long sendTimeout) {
|
||||
this.sendTimeout = sendTimeout;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Object doHandle(Message message, SimpleMethodInvoker invoker) {
|
||||
if (method.getParameterTypes().length != 1) {
|
||||
@@ -139,7 +145,7 @@ public class SplitterMessageHandlerAdapter extends AbstractMessageHandlerAdapter
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("sending message to channel '" + channelName + "'");
|
||||
}
|
||||
return channel.send(message);
|
||||
return (this.sendTimeout < 0) ? channel.send(message) : channel.send(message, this.sendTimeout);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -76,9 +76,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(3, dispatcherTask.dispatch());
|
||||
byte[] result = stream.toByteArray();
|
||||
assertEquals(9, result.length);
|
||||
@@ -95,9 +95,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
byte[] result = stream.toByteArray();
|
||||
assertEquals(6, result.length);
|
||||
@@ -114,9 +114,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(3, dispatcherTask.dispatch());
|
||||
byte[] result = stream.toByteArray();
|
||||
assertEquals(9, result.length);
|
||||
@@ -133,9 +133,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
byte[] result1 = stream.toByteArray();
|
||||
assertEquals(6, result1.length);
|
||||
@@ -157,9 +157,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(3, dispatcherTask.dispatch());
|
||||
byte[] result1 = stream.toByteArray();
|
||||
assertEquals(9, result1.length);
|
||||
@@ -180,9 +180,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
byte[] result1 = stream.toByteArray();
|
||||
assertEquals(6, result1.length);
|
||||
@@ -203,9 +203,9 @@ public class ByteStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}));
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
byte[] result1 = stream.toByteArray();
|
||||
assertEquals(6, result1.length);
|
||||
|
||||
@@ -55,8 +55,8 @@ public class CharacterStreamTargetAdapterTests {
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new StringMessage("foo"));
|
||||
channel.send(new StringMessage("bar"));
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
assertEquals(1, dispatcherTask.dispatch());
|
||||
String result1 = new String(stream.toByteArray());
|
||||
assertEquals("foo", result1);
|
||||
@@ -73,8 +73,8 @@ public class CharacterStreamTargetAdapterTests {
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new StringMessage("foo"));
|
||||
channel.send(new StringMessage("bar"));
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
assertEquals(1, dispatcherTask.dispatch());
|
||||
String result1 = new String(stream.toByteArray());
|
||||
String newLine = System.getProperty("line.separator");
|
||||
@@ -93,8 +93,8 @@ public class CharacterStreamTargetAdapterTests {
|
||||
SimpleChannel channel = new SimpleChannel(dispatcherPolicy);
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new StringMessage("foo"));
|
||||
channel.send(new StringMessage("bar"));
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
String result = new String(stream.toByteArray());
|
||||
assertEquals("foobar", result);
|
||||
@@ -111,8 +111,8 @@ public class CharacterStreamTargetAdapterTests {
|
||||
DispatcherTask dispatcherTask = new DispatcherTask(channel);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
dispatcherTask.addHandler(adapter);
|
||||
channel.send(new StringMessage("foo"));
|
||||
channel.send(new StringMessage("bar"));
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
String result = new String(stream.toByteArray());
|
||||
String newLine = System.getProperty("line.separator");
|
||||
@@ -146,8 +146,8 @@ public class CharacterStreamTargetAdapterTests {
|
||||
dispatcherTask.addHandler(adapter);
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
TestObject testObject2 = new TestObject("bar");
|
||||
channel.send(new GenericMessage<TestObject>(testObject1));
|
||||
channel.send(new GenericMessage<TestObject>(testObject2));
|
||||
channel.send(new GenericMessage<TestObject>(testObject1), 0);
|
||||
channel.send(new GenericMessage<TestObject>(testObject2), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
String result = new String(stream.toByteArray());
|
||||
assertEquals("foobar", result);
|
||||
@@ -166,8 +166,8 @@ public class CharacterStreamTargetAdapterTests {
|
||||
dispatcherTask.addHandler(adapter);
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
TestObject testObject2 = new TestObject("bar");
|
||||
channel.send(new GenericMessage<TestObject>(testObject1));
|
||||
channel.send(new GenericMessage<TestObject>(testObject2));
|
||||
channel.send(new GenericMessage<TestObject>(testObject1), 0);
|
||||
channel.send(new GenericMessage<TestObject>(testObject2), 0);
|
||||
assertEquals(2, dispatcherTask.dispatch());
|
||||
String result = new String(stream.toByteArray());
|
||||
String newLine = System.getProperty("line.separator");
|
||||
|
||||
@@ -42,8 +42,8 @@ public class ChannelPollingMessageRetrieverTests {
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel);
|
||||
Collection<Message<?>> results = retriever.retrieveMessages();
|
||||
assertTrue(results.isEmpty());
|
||||
channel.send(new StringMessage("test1"));
|
||||
channel.send(new StringMessage("test2"));
|
||||
channel.send(new StringMessage("test1"), 0);
|
||||
channel.send(new StringMessage("test2"), 0);
|
||||
results = retriever.retrieveMessages();
|
||||
assertEquals(1, results.size());
|
||||
assertEquals("test1", results.iterator().next().getPayload());
|
||||
@@ -61,9 +61,9 @@ public class ChannelPollingMessageRetrieverTests {
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel);
|
||||
Collection<Message<?>> results = retriever.retrieveMessages();
|
||||
assertTrue(results.isEmpty());
|
||||
channel.send(new StringMessage("test1"));
|
||||
channel.send(new StringMessage("test2"));
|
||||
channel.send(new StringMessage("test3"));
|
||||
channel.send(new StringMessage("test1"), 0);
|
||||
channel.send(new StringMessage("test2"), 0);
|
||||
channel.send(new StringMessage("test3"), 0);
|
||||
results = retriever.retrieveMessages();
|
||||
assertEquals(2, results.size());
|
||||
Iterator<Message<?>> iter = results.iterator();
|
||||
@@ -83,9 +83,9 @@ public class ChannelPollingMessageRetrieverTests {
|
||||
ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel);
|
||||
Collection<Message<?>> results = retriever.retrieveMessages();
|
||||
assertTrue(results.isEmpty());
|
||||
channel.send(new StringMessage("test1"));
|
||||
channel.send(new StringMessage("test2"));
|
||||
channel.send(new StringMessage("test3"));
|
||||
channel.send(new StringMessage("test1"), 0);
|
||||
channel.send(new StringMessage("test2"), 0);
|
||||
channel.send(new StringMessage("test3"), 0);
|
||||
results = retriever.retrieveMessages();
|
||||
assertEquals(1, results.size());
|
||||
assertEquals("test1", results.iterator().next().getPayload());
|
||||
|
||||
Reference in New Issue
Block a user