Refactored CharacterStreamOutboundChannelAdapter to CharacterStreamWritingMessageConsumer and simplified the abstract method for AbstractOutboundChannelAdapter so that only a bean definition is returned (the base class now handles registration).

This commit is contained in:
Mark Fisher
2008-09-24 00:00:55 +00:00
parent db33965e77
commit 0ed3ba7657
9 changed files with 116 additions and 131 deletions

View File

@@ -27,21 +27,21 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.integration.ConfigurationException;
import org.springframework.integration.endpoint.AbstractMessageConsumingEndpoint;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageConsumer;
import org.springframework.integration.message.MessagingException;
import org.springframework.util.Assert;
/**
* An outbound Channel Adapter that writes to a {@link Writer}. String-based objects
* will be written directly, but if the object is not itself a {@link String}, the
* target will write the result of the object's {@link #toString()} method. To
* append a new-line after each write, set the {@link #shouldAppendNewLine} flag to
* <em>true</em>. It is <em>false</em> by default.
* A {@link MessageConsumer} that writes characters to a {@link Writer}.
* String, character array, and byte array payloads will be written directly,
* but for other payload types, the result of the object's {@link #toString()}
* method will be written. To append a new-line after each write, set the
* {@link #shouldAppendNewLine} flag to 'true'. It is 'false' by default.
*
* @author Mark Fisher
*/
public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsumingEndpoint {
public class CharacterStreamWritingMessageConsumer implements MessageConsumer {
private final Log logger = LogFactory.getLog(this.getClass());
@@ -50,11 +50,11 @@ public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsum
private volatile boolean shouldAppendNewLine = false;
public CharacterStreamOutboundChannelAdapter(Writer writer) {
public CharacterStreamWritingMessageConsumer(Writer writer) {
this(writer, -1);
}
public CharacterStreamOutboundChannelAdapter(Writer writer, int bufferSize) {
public CharacterStreamWritingMessageConsumer(Writer writer, int bufferSize) {
Assert.notNull(writer, "writer must not be null");
if (writer instanceof BufferedWriter) {
this.writer = (BufferedWriter) writer;
@@ -72,7 +72,7 @@ public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsum
* Factory method that creates a target for stdout (System.out) with the
* default charset encoding.
*/
public static CharacterStreamOutboundChannelAdapter stdout() {
public static CharacterStreamWritingMessageConsumer stdout() {
return stdout(null);
}
@@ -80,7 +80,7 @@ public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsum
* Factory method that creates a target for stdout (System.out) with the
* specified charset encoding.
*/
public static CharacterStreamOutboundChannelAdapter stdout(String charsetName) {
public static CharacterStreamWritingMessageConsumer stdout(String charsetName) {
return createTargetForStream(System.out, charsetName);
}
@@ -88,7 +88,7 @@ public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsum
* Factory method that creates a target for stderr (System.err) with the
* default charset encoding.
*/
public static CharacterStreamOutboundChannelAdapter stderr() {
public static CharacterStreamWritingMessageConsumer stderr() {
return stderr(null);
}
@@ -96,16 +96,16 @@ public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsum
* Factory method that creates a target for stderr (System.err) with the
* specified charset encoding.
*/
public static CharacterStreamOutboundChannelAdapter stderr(String charsetName) {
public static CharacterStreamWritingMessageConsumer stderr(String charsetName) {
return createTargetForStream(System.err, charsetName);
}
private static CharacterStreamOutboundChannelAdapter createTargetForStream(OutputStream stream, String charsetName) {
private static CharacterStreamWritingMessageConsumer createTargetForStream(OutputStream stream, String charsetName) {
if (charsetName == null) {
return new CharacterStreamOutboundChannelAdapter(new OutputStreamWriter(stream));
return new CharacterStreamWritingMessageConsumer(new OutputStreamWriter(stream));
}
try {
return new CharacterStreamOutboundChannelAdapter(new OutputStreamWriter(stream, charsetName));
return new CharacterStreamWritingMessageConsumer(new OutputStreamWriter(stream, charsetName));
}
catch (UnsupportedEncodingException e) {
throw new ConfigurationException("unsupported encoding: " + charsetName, e);
@@ -117,8 +117,7 @@ public class CharacterStreamOutboundChannelAdapter extends AbstractMessageConsum
this.shouldAppendNewLine = shouldAppendNewLine;
}
@Override
public void onMessageInternal(Message<?> message) {
public void onMessage(Message<?> message) {
Object payload = message.getPayload();
if (payload == null) {
if (logger.isWarnEnabled()) {

View File

@@ -18,16 +18,11 @@ package org.springframework.integration.stream.config;
import org.w3c.dom.Element;
import org.springframework.beans.factory.BeanDefinitionStoreException;
import org.springframework.beans.factory.config.BeanDefinitionHolder;
import org.springframework.beans.factory.support.AbstractBeanDefinition;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.support.BeanDefinitionReaderUtils;
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
import org.springframework.beans.factory.xml.ParserContext;
import org.springframework.integration.ConfigurationException;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.stream.CharacterStreamOutboundChannelAdapter;
import org.springframework.integration.config.AbstractOutboundChannelAdapterParser;
import org.springframework.integration.stream.CharacterStreamWritingMessageConsumer;
import org.springframework.util.StringUtils;
/**
@@ -35,28 +30,12 @@ import org.springframework.util.StringUtils;
*
* @author Mark Fisher
*/
public class ConsoleOutboundChannelAdapterParser extends AbstractSingleBeanDefinitionParser {
public class ConsoleOutboundChannelAdapterParser extends AbstractOutboundChannelAdapterParser {
@Override
protected Class<?> getBeanClass(Element element) {
return CharacterStreamOutboundChannelAdapter.class;
}
@Override
protected String resolveId(Element element, AbstractBeanDefinition definition, ParserContext parserContext) throws BeanDefinitionStoreException {
String id = element.getAttribute("id");
if (!element.hasAttribute("channel")) {
// the created channel will get the 'id', so the adapter's bean name includes a suffix
id = id + ".adapter";
}
else if (!StringUtils.hasText(id)) {
id = parserContext.getReaderContext().generateBeanName(definition);
}
return id;
}
@Override
protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
protected AbstractBeanDefinition parseConsumer(Element element, ParserContext parserContext) {
BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(
CharacterStreamWritingMessageConsumer.class);
if (element.getLocalName().startsWith("stderr")) {
builder.setFactoryMethod("stderr");
}
@@ -70,25 +49,7 @@ public class ConsoleOutboundChannelAdapterParser extends AbstractSingleBeanDefin
if ("true".equals(element.getAttribute("append-newline"))) {
builder.addPropertyValue("shouldAppendNewLine", Boolean.TRUE);
}
String channelName = element.getAttribute("channel");
if (StringUtils.hasText(channelName)) {
builder.addPropertyReference("inputChannel", channelName);
}
else {
builder.addPropertyReference("inputChannel", this.createDirectChannel(element, parserContext));
}
}
private String createDirectChannel(Element element, ParserContext parserContext) {
String channelId = element.getAttribute("id");
if (!StringUtils.hasText(channelId)) {
throw new ConfigurationException("The channel-adapter's 'id' attribute is required when no 'channel' "
+ "reference has been provided, because that 'id' would be used for the created channel.");
}
BeanDefinitionBuilder channelBuilder = BeanDefinitionBuilder.genericBeanDefinition(DirectChannel.class);
BeanDefinitionHolder holder = new BeanDefinitionHolder(channelBuilder.getBeanDefinition(), channelId);
BeanDefinitionReaderUtils.registerBeanDefinition(holder, parserContext.getRegistry());
return channelId;
return builder.getBeanDefinition();
}
}

View File

@@ -32,7 +32,7 @@ import org.springframework.integration.scheduling.PollingSchedule;
/**
* @author Mark Fisher
*/
public class CharacterStreamOutboundChannelAdapterTests {
public class CharacterStreamWritingMessageConsumerTests {
private QueueChannel channel;
@@ -47,18 +47,18 @@ public class CharacterStreamOutboundChannelAdapterTests {
@Test
public void testSingleString() {
public void singleString() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
target.onMessage(new StringMessage("foo"));
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
consumer.onMessage(new StringMessage("foo"));
assertEquals("foo", writer.toString());
}
@Test
public void testTwoStringsAndNoNewLinesByDefault() {
public void twoStringsAndNoNewLinesByDefault() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
poller.subscribe(target);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
poller.subscribe(consumer);
poller.setMaxMessagesPerPoll(1);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
@@ -69,11 +69,11 @@ public class CharacterStreamOutboundChannelAdapterTests {
}
@Test
public void testTwoStringsWithNewLines() {
public void twoStringsWithNewLines() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
target.setShouldAppendNewLine(true);
poller.subscribe(target);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
consumer.setShouldAppendNewLine(true);
poller.subscribe(consumer);
poller.setMaxMessagesPerPoll(1);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
@@ -85,11 +85,11 @@ public class CharacterStreamOutboundChannelAdapterTests {
}
@Test
public void testMaxMessagesPerTaskSameAsMessageCount() {
public void maxMessagesPerTaskSameAsMessageCount() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
poller.setMaxMessagesPerPoll(2);
poller.subscribe(target);
poller.subscribe(consumer);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
poller.run();
@@ -97,13 +97,13 @@ public class CharacterStreamOutboundChannelAdapterTests {
}
@Test
public void testMaxMessagesPerTaskExceedsMessageCountWithAppendedNewLines() {
public void maxMessagesPerTaskExceedsMessageCountWithAppendedNewLines() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
poller.setMaxMessagesPerPoll(10);
poller.setReceiveTimeout(0);
poller.subscribe(target);
target.setShouldAppendNewLine(true);
poller.subscribe(consumer);
consumer.setShouldAppendNewLine(true);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
poller.run();
@@ -112,10 +112,10 @@ public class CharacterStreamOutboundChannelAdapterTests {
}
@Test
public void testSingleNonStringObject() {
public void singleNonStringObject() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
poller.subscribe(target);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
poller.subscribe(consumer);
poller.setMaxMessagesPerPoll(1);
TestObject testObject = new TestObject("foo");
channel.send(new GenericMessage<TestObject>(testObject));
@@ -124,12 +124,12 @@ public class CharacterStreamOutboundChannelAdapterTests {
}
@Test
public void testTwoNonStringObjectWithOutNewLines() {
public void twoNonStringObjectWithOutNewLines() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
poller.setReceiveTimeout(0);
poller.setMaxMessagesPerPoll(2);
poller.subscribe(target);
poller.subscribe(consumer);
TestObject testObject1 = new TestObject("foo");
TestObject testObject2 = new TestObject("bar");
channel.send(new GenericMessage<TestObject>(testObject1), 0);
@@ -139,13 +139,13 @@ public class CharacterStreamOutboundChannelAdapterTests {
}
@Test
public void testTwoNonStringObjectWithNewLines() {
public void twoNonStringObjectWithNewLines() {
StringWriter writer = new StringWriter();
CharacterStreamOutboundChannelAdapter target = new CharacterStreamOutboundChannelAdapter(writer);
target.setShouldAppendNewLine(true);
CharacterStreamWritingMessageConsumer consumer = new CharacterStreamWritingMessageConsumer(writer);
consumer.setShouldAppendNewLine(true);
poller.setReceiveTimeout(0);
poller.setMaxMessagesPerPoll(2);
poller.subscribe(target);
poller.subscribe(consumer);
TestObject testObject1 = new TestObject("foo");
TestObject testObject2 = new TestObject("bar");
channel.send(new GenericMessage<TestObject>(testObject1), 0);

View File

@@ -34,7 +34,7 @@ import org.springframework.beans.factory.BeanCreationException;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.ConfigurationException;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.stream.CharacterStreamOutboundChannelAdapter;
import org.springframework.integration.stream.CharacterStreamWritingMessageConsumer;
/**
* @author Mark Fisher
@@ -61,9 +61,10 @@ public class ConsoleOutboundChannelAdapterParserTests {
public void stdoutAdapterWithDefaultCharset() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"consoleOutboundChannelAdapterParserTests.xml", ConsoleOutboundChannelAdapterParserTests.class);
CharacterStreamOutboundChannelAdapter adapter =
(CharacterStreamOutboundChannelAdapter) context.getBean("stdoutAdapterWithDefaultCharset");
DirectFieldAccessor accessor = new DirectFieldAccessor(adapter);
Object adapter = context.getBean("stdoutAdapterWithDefaultCharset");
CharacterStreamWritingMessageConsumer consumer = (CharacterStreamWritingMessageConsumer)
new DirectFieldAccessor(adapter).getPropertyValue("consumer");
DirectFieldAccessor accessor = new DirectFieldAccessor(consumer);
Writer bufferedWriter = (Writer) accessor.getPropertyValue("writer");
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
DirectFieldAccessor bufferedWriterAccessor = new DirectFieldAccessor(bufferedWriter);
@@ -72,7 +73,7 @@ public class ConsoleOutboundChannelAdapterParserTests {
Charset writerCharset = Charset.forName(((OutputStreamWriter) writer).getEncoding());
assertEquals(Charset.defaultCharset(), writerCharset);
this.resetStreams();
adapter.onMessage(new StringMessage("foo"));
consumer.onMessage(new StringMessage("foo"));
assertEquals("foo", out.toString());
assertEquals("", err.toString());
}
@@ -81,9 +82,10 @@ public class ConsoleOutboundChannelAdapterParserTests {
public void stdoutAdapterWithProvidedCharset() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"consoleOutboundChannelAdapterParserTests.xml", ConsoleOutboundChannelAdapterParserTests.class);
CharacterStreamOutboundChannelAdapter adapter =
(CharacterStreamOutboundChannelAdapter) context.getBean("stdoutAdapterWithProvidedCharset");
DirectFieldAccessor accessor = new DirectFieldAccessor(adapter);
Object adapter = context.getBean("stdoutAdapterWithProvidedCharset");
CharacterStreamWritingMessageConsumer consumer = (CharacterStreamWritingMessageConsumer)
new DirectFieldAccessor(adapter).getPropertyValue("consumer");
DirectFieldAccessor accessor = new DirectFieldAccessor(consumer);
Writer bufferedWriter = (Writer) accessor.getPropertyValue("writer");
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
DirectFieldAccessor bufferedWriterAccessor = new DirectFieldAccessor(bufferedWriter);
@@ -92,7 +94,7 @@ public class ConsoleOutboundChannelAdapterParserTests {
Charset writerCharset = Charset.forName(((OutputStreamWriter) writer).getEncoding());
assertEquals(Charset.forName("UTF-8"), writerCharset);
this.resetStreams();
adapter.onMessage(new StringMessage("bar"));
consumer.onMessage(new StringMessage("bar"));
assertEquals("bar", out.toString());
assertEquals("", err.toString());
}
@@ -117,9 +119,10 @@ public class ConsoleOutboundChannelAdapterParserTests {
public void stderrAdapter() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"consoleOutboundChannelAdapterParserTests.xml", ConsoleOutboundChannelAdapterParserTests.class);
CharacterStreamOutboundChannelAdapter adapter =
(CharacterStreamOutboundChannelAdapter) context.getBean("stderrAdapter");
DirectFieldAccessor accessor = new DirectFieldAccessor(adapter);
Object adapter = context.getBean("stderrAdapter");
CharacterStreamWritingMessageConsumer consumer = (CharacterStreamWritingMessageConsumer)
new DirectFieldAccessor(adapter).getPropertyValue("consumer");
DirectFieldAccessor accessor = new DirectFieldAccessor(consumer);
Writer bufferedWriter = (Writer) accessor.getPropertyValue("writer");
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
DirectFieldAccessor bufferedWriterAccessor = new DirectFieldAccessor(bufferedWriter);
@@ -128,7 +131,7 @@ public class ConsoleOutboundChannelAdapterParserTests {
Charset writerCharset = Charset.forName(((OutputStreamWriter) writer).getEncoding());
assertEquals(Charset.defaultCharset(), writerCharset);
this.resetStreams();
adapter.onMessage(new StringMessage("bad"));
consumer.onMessage(new StringMessage("bad"));
assertEquals("", out.toString());
assertEquals("bad", err.toString());
}
@@ -137,9 +140,10 @@ public class ConsoleOutboundChannelAdapterParserTests {
public void stdoutAdatperWithAppendNewLine() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"consoleOutboundChannelAdapterParserTests.xml", ConsoleOutboundChannelAdapterParserTests.class);
CharacterStreamOutboundChannelAdapter adapter =
(CharacterStreamOutboundChannelAdapter) context.getBean("newlineAdapter");
DirectFieldAccessor accessor = new DirectFieldAccessor(adapter);
Object adapter = context.getBean("newlineAdapter");
CharacterStreamWritingMessageConsumer consumer = (CharacterStreamWritingMessageConsumer)
new DirectFieldAccessor(adapter).getPropertyValue("consumer");
DirectFieldAccessor accessor = new DirectFieldAccessor(consumer);
Writer bufferedWriter = (Writer) accessor.getPropertyValue("writer");
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
DirectFieldAccessor bufferedWriterAccessor = new DirectFieldAccessor(bufferedWriter);
@@ -148,7 +152,7 @@ public class ConsoleOutboundChannelAdapterParserTests {
Charset writerCharset = Charset.forName(((OutputStreamWriter) writer).getEncoding());
assertEquals(Charset.defaultCharset(), writerCharset);
this.resetStreams();
adapter.onMessage(new StringMessage("foo"));
consumer.onMessage(new StringMessage("foo"));
assertEquals("foo\n", out.toString());
}