This commit is contained in:
Mark Fisher
2009-07-03 19:52:23 +00:00
parent bd81a6fbe2
commit 6315a85b62

View File

@@ -49,8 +49,8 @@
<title>SubscribableChannel</title>
<para>
The <interfacename>SubscribableChannel</interfacename> base interface is implemented by channels that send
Messages directly to their subscribed handlers. Therefore, they do not provide receive methods for polling, but
instead define methods for handling those subscribers:
Messages directly to their subscribed <interfacename>MessageHandler</interfacename>s. Therefore, they do not
provide receive methods for polling, but instead define methods for managing those subscribers:
<programlisting language="java">public interface SubscribableChannel extends MessageChannel {
boolean subscribe(MessageHandler handler);
@@ -95,10 +95,12 @@
<programlisting language="java">public QueueChannel(int capacity)</programlisting>
A channel that has not reached its capacity limit will store messages in its internal queue, and the
<methodname>send()</methodname> method will return immediately even if no receiver is ready to handle the
message. If the queue has reached capacity, then the sender will block until room is available. Likewise, a
receive call will return immediately if a message is available on the queue, but if the queue is empty, then
a receive call may block until either a message is available or the timeout elapses. In either case, it is
possible to force an immediate return regardless of the queue's state by passing a timeout value of 0.
message. If the queue has reached capacity, then the sender will block until room is available. Or, if using
the send call that accepts a timeout, it will block until either room is available or the timeout period
elapses, whichever occurs first. Likewise, a receive call will return immediately if a message is available
on the queue, but if the queue is empty, then a receive call may block until either a message is available
or the timeout elapses. In either case, it is possible to force an immediate return regardless of the
queue's state by passing a timeout value of 0. The no-argument send and receive methods block indefinitely.
Note however, that calling the no-arg versions of <methodname>send()</methodname> and
<methodname>receive()</methodname> will block indefinitely.
</para>
@@ -107,11 +109,11 @@
<title>PriorityChannel</title>
<para>
Whereas the <classname>QueueChannel</classname> enforces first-in/first-out (FIFO) ordering, the
<classname>PriorityChannel</classname> is an alternative implementation that allows for messages to be ordered
within the channel based upon a priority. By default the priority is determined by the
<classname>PriorityChannel</classname> is an alternative implementation that allows for messages
to be ordered within the channel based upon a priority. By default the priority is determined by the
'<literal>priority</literal>' header within each message. However, for custom priority determination
logic, a comparator of type <classname>Comparator&lt;Message&lt;?&gt;&gt;</classname> can be provided to the
<classname>PriorityChannel</classname>'s constructor.
logic, a comparator of type <classname>Comparator&lt;Message&lt;?&gt;&gt;</classname> can be provided
to the <classname>PriorityChannel</classname>'s constructor.
</para>
</section>
<section id="channel-implementations-rendezvouschannel">
@@ -122,47 +124,90 @@
this implementation is quite similar to the <classname>QueueChannel</classname> except that it uses a
<classname>SynchronousQueue</classname> (a zero-capacity implementation of
<interfacename>BlockingQueue</interfacename>). This works well in situations where the sender and receiver are
operating in different threads but simply dropping the message in a queue asynchronously is too dangerous. For
example, the sender's thread could roll back a transaction if the send operation times out, whereas with a
<classname>QueueChannel</classname>, the message would have been stored to the internal queue and potentially
never received.
operating in different threads but simply dropping the message in a queue asynchronously is not appropriate.
In other words, with a <classname>RendezvousChannel</classname> at least the sender knows that some receiver
has accepted the message, whereas with a <classname>QueueChannel</classname>, the message would have been
stored to the internal queue and potentially never received.
</para>
<tip>
<para>
Keep in mind that all of these queue-based channels are storing messages in-memory only. When persistence
is required, you can either invoke a database operation within a handler or use Spring Integration's
support for JMS-based Channel Adapters. The latter option allows you to take advantage of any JMS provider's
implementation for message persistence, and it will be discussed in <xref linkend="jms"/>. However, when
buffering in a queue is not necessary, the simplest approach is to rely upon the
<classname>DirectChannel</classname> discussed next.
</para>
</tip>
<para>
The <classname>RendezvousChannel</classname> is also useful for implementing request-reply
operations. The sender can create a temporary, anonymous instance of <classname>RendezvousChannel</classname>
which it then sets as the 'replyChannel' header when building a Message. After sending that Message, the sender
can immediately call receive (optionally providing a timeout value) in order to block while waiting for a reply
Message.
Message. This is very similar to the implementation used internally by many of Spring Integration's
request-reply components.
</para>
</section>
<section id="channel-implementations-directchannel">
<title>DirectChannel</title>
<para>
The <classname>DirectChannel</classname> has point-to-point semantics, but otherwise is more similar to the
The <classname>DirectChannel</classname> has point-to-point semantics but otherwise is more similar to the
<classname>PublishSubscribeChannel</classname> than any of the queue-based channel implementations described
above. It implements the <interfacename>SubscribableChannel</interfacename> interface instead of the
<interfacename>PollableChannel</interfacename> interface, so it dispatches Messages directly to a subscriber.
As a point-to-point channel, however, it differs from the <classname>PublishSubscribeChannel</classname> in
that it will only send each Message to a <emphasis>single</emphasis> subscribed
<classname>MessageHandler</classname>. Its primary purpose is to enable a single thread to perform the
operations on "both sides" of the channel. For example, if a handler is subscribed to a
<classname>DirectChannel</classname>, then sending a Message to that channel will trigger invocation of that
handler's <methodname>handleMessage(Message)</methodname> method <emphasis>directly in the sender's
thread</emphasis>. The key motivation for providing a channel implementation with this behavior is to support
transactions that must span across the channel while still benefiting from the abstraction and loose coupling
that the channel provides. If the send call is invoked within the scope of a transaction, then the outcome of
the handler's invocation (e.g. updating a database record) can play a role in determining the ultimate result
of that transaction (commit or rollback).
<classname>MessageHandler</classname>.
</para>
<para>
In addition to being the simplest point-to-point channel option, one of its most important features is that
it enables a single thread to perform the operations on "both sides" of the channel. For example, if a handler
is subscribed to a <classname>DirectChannel</classname>, then sending a Message to that channel will trigger
invocation of that handler's <methodname>handleMessage(Message)</methodname> method <emphasis>directly in the
sender's thread</emphasis>, before the send() method invocation can return.
</para>
<para>
The key motivation for providing a channel implementation with this behavior is to support transactions that
must span across the channel while still benefiting from the abstraction and loose coupling that the channel
provides. If the send call is invoked within the scope of a transaction, then the outcome of the handler's
invocation (e.g. updating a database record) will play a role in determining the ultimate result of that
transaction (commit or rollback).
<note>
Since the <classname>DirectChannel</classname> is the simplest option and does not add any additional
overhead that would be required for scheduling and managing the threads of a poller, it is the default
channel type within Spring Integration. The general idea is to define the channels for an application and
then to consider which of those needs to provide buffering to throttle input, and to modify those to be
queue-based <interfacename>PollableChannels</interfacename>. Likewise, if a channel needs to broadcast
then to consider which of those need to provide buffering or to throttle input, and then modify those to
be queue-based <interfacename>PollableChannels</interfacename>. Likewise, if a channel needs to broadcast
messages, it should not be a <classname>DirectChannel</classname> but rather a
<classname>PublishSubscribeChannel</classname>. Below you will see how each of these can be configured.
</note>
</para>
<para>
The <classname>DirectChannel</classname> can have one of two dispatcher strategies. These determine how
invocations will be ordered in the case that there are multiple handlers subscribed to the same channel.
The default strategy is "round-robin" and essentially load-balances across the handlers in rotation. The
other strategy is "failover" and it will always try to invoke the first handler, falling back to any
subsequent handlers as necessary. The order is determined by an optional order value defined on the
handlers themselves or, if no such value exists, the order in which the handlers are subscribed.
</para>
</section>
<section id="executor-channel">
<title>ExecutorChannel</title>
<para>
The <classname>ExecutorChannel</classname> is a point-to-point channel that supports
the same dispatcher strategies as <classname>DirectChannel</classname>. The key difference is that
it delegates to an instance of <interfacename>TaskExecutor</interfacename> to perform the dispatch.
This means that the send method typically will not block, but it also means that the handler
invocation may not occur in the sender's thread. It therefore <emphasis>does not support
transactions spanning the sender and receiving handler</emphasis>.
<tip>
Note that there are occasions where the sender may block. For example, when using a
TaskExecutor with a rejection-policy that throttles back on the client (such as the
<code>ThreadPoolExecutor.CallerRunsPolicy</code>), the sender's thread will execute
the method directly anytime the thread pool is at its maximum capacity and the
executor's work queue is full.
</tip>
</para>
</section>
<section id="channel-implementations-threadlocalchannel">
<title>ThreadLocalChannel</title>
@@ -200,15 +245,25 @@
After implementing the interface, registering the interceptor with a channel is just a matter of calling:
<programlisting language="java">channel.addInterceptor(someChannelInterceptor);</programlisting>
The methods that return a Message instance can be used for transforming the Message or can return 'null'
to prevent further processing (of course, any of the methods can throw an Exception). Also, the
to prevent further processing (of course, any of the methods can throw a RuntimeException). Also, the
<methodname>preReceive</methodname> method can return '<literal>false</literal>' to prevent the receive
operation from proceeding.
<note>
Keep in mind that <methodname>receive()</methodname> calls are only relevant for
<interfacename>PollableChannels</interfacename>. In fact the
<interfacename>SubscribableChannel</interfacename> interface does not even define a
<methodname>receive()</methodname> method. The reason for this is that when a Message is sent to a
<interfacename>SubscribableChannel</interfacename> it will be sent directly to one or more subscribers
depending on the type of channel (e.g. a PublishSubscribeChannel sends to all of its subscribers). Therefore,
the <methodname>preReceive(..)</methodname> and <methodname>postReceive(..)</methodname> interceptor methods
are only invoked when the interceptor is applied to a <interfacename>PollableChannel</interfacename>.
</note>
</para>
<para>
Because it is rarely necessary to implement all of the interceptor methods, a
<classname>ChannelInterceptorAdapter</classname> class is also available for sub-classing. It provides no-op
methods (the <literal>void</literal> method is empty, the <classname>Message</classname> returning methods
return the Message parameter as-is, and the <literal>boolean</literal> method returns <literal>true</literal>).
return the Message as-is, and the <literal>boolean</literal> method returns <literal>true</literal>).
Therefore, it is often easiest to extend that class and just implement the method(s) that you need as in the
following example.
<programlisting language="java"><![CDATA[public class CountingChannelInterceptor extends ChannelInterceptorAdapter {
@@ -221,16 +276,19 @@
return message;
}
}]]></programlisting>
<note>
Keep in mind that <methodname>receive()</methodname> calls are only relevant for
<interfacename>PollableChannels</interfacename>. In fact the
<interfacename>SubscribableChannel</interfacename> interface does not even define a
<methodname>receive()</methodname> method. The reason for this is that when a Message is sent to a
<interfacename>SubscribableChannel</interfacename> it will be sent directly to one or more subscribers
depending on the type of channel (e.g. a PublishSubscribeChannel sends to all of its subscribers). Therefore,
the <methodname>preReceive(..)</methodname> and <methodname>postReceive(..)</methodname> interceptor methods
are only invoked when the interceptor is applied to a <interfacename>PollableChannel</interfacename>.
</note>
<tip>
The order of invocation for the interceptor methods depends on the type of channel. As described above,
the queue-based channels are the only ones where the receive method is intercepted in the first place.
Additionally, the relationship between send and receive interception depends on the timing of separate
sender and receiver threads. For example, if a receiver is already blocked while waiting for a message
the order could be: preSend, preReceive, postReceive, postSend. However, if a receiver polls after the
sender has placed a message on the channel and already returned, the order would be: preSend, postSend,
(some-time-elapses) preReceive, postReceive. The time that elapses in such a case depends on a number
of factors and is therefore generally unpredictable (in fact, the receive may never happen!).
Obviously, the type of queue also plays a role (e.g. rendezvous vs. priority). The bottom line is that
you cannot rely on the order beyond the fact that preSend will precede postSend and preReceive will
precede postReceive.
</tip>
</para>
</section>
@@ -256,6 +314,12 @@ public Message<?> sendAndReceive(final Message<?> request, final MessageChannel
public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programlisting>
</para>
<note>
<para>
A less invasive approach that allows you to invoke simple interfaces with payload and/or header
values instead of Message instances is described in <xref linkend="gateway-proxy"/>.
</para>
</note>
</section>
<section id="channel-configuration">
@@ -285,14 +349,15 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
instance (a <interfacename>SubscribableChannel</interfacename>).
</para>
<para>
However, you can also provide a variety of "queue" sub-elements to create the channel types (as described in
However, you can alternatively provide a variety of "queue" sub-elements to create any of
the pollable channel types (as described in
<xref linkend="channel-implementations"/>). Examples of each are shown below.
</para>
<section id="channel-configuration-directchannel">
<title>DirectChannel Configuration</title>
<para>
As mentioned above, <classname>DirectChannel</classname> is the default type.
<programlisting language="xml"><![CDATA[<channel id="exampleChannel"/>]]></programlisting>
As mentioned above, <classname>DirectChannel</classname> is the default type.
<programlisting language="xml"><![CDATA[<channel id="directChannel"/>]]></programlisting>
</para>
</section>
<section id="channel-configuration-queuechannel">
@@ -300,12 +365,12 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
<para>
To create a <classname>QueueChannel</classname>, use the "queue" sub-element.
You may specify the channel's capacity:
<programlisting language="xml">&lt;channel id="exampleChannel"&gt;
<programlisting language="xml">&lt;channel id="queueChannel"&gt;
&lt;queue capacity="25"/&gt;
&lt;/channel&gt;</programlisting>
<note>
If you do not provide a value for the 'capacity' attribute on this &lt;queue/&gt; sub-element,
the resulting queue will be unbounded. To avoid issues such as OutOfMemoryErrors, it's highly
the resulting queue will be unbounded. To avoid issues such as OutOfMemoryErrors, it is highly
recommended to set an explicit value for a bounded queue.
</note>
</para>
@@ -316,21 +381,41 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
To create a <classname>PublishSubscribeChannel</classname>, use the "publish-subscribe-channel" element.
When using this element, you can also specify the "task-executor" used for publishing
Messages (if none is specified it simply publishes in the sender's thread):
<programlisting language="xml">&lt;publish-subscribe-channel id="exampleChannel" task-executor="someTaskExecutor"/&gt;</programlisting>
<programlisting language="xml">&lt;publish-subscribe-channel id="pubsubChannel" task-executor="someExecutor"/&gt;</programlisting>
If you are providing a <emphasis>Resequencer</emphasis> or <emphasis>Aggregator</emphasis> downstream
from a <classname>PublishSubscribeChannel</classname>, then you can set the 'apply-sequence' property
for the channel. That will indicate that the channel should set the sequence-size and sequence-number
Message headers prior to passing the Messages along. For example, if there are 5 subscribers, the
sequence-size would be set to 5, and the Messages would have sequence-number header values ranging
from 1 to 5. This value is 'false' by default.
<programlisting language="xml">&lt;publish-subscribe-channel id="exampleChannel" apply-sequence="true"/&gt;</programlisting>
on the channel to <code>true</code>. That will indicate that the channel should set the sequence-size
and sequence-number Message headers as well as the correlation id prior to passing the Messages along.
For example, if there are 5 subscribers, the sequence-size would be set to 5, and the Messages would
have sequence-number header values ranging from 1 to 5.
<programlisting language="xml">&lt;publish-subscribe-channel id="pubsubChannel" apply-sequence="true"/&gt;</programlisting>
<note>
The 'apply-sequence' value is <code>false</code> by default so that a Publish Subscribe Channel
can send the exact same Message instances to multiple outbound channels. Since Spring Integration
enforces immutability of the payload and header references, the channel creates new Message
instances with the same payload reference but different header values when the flag is set to
<code>true</code>.
</note>
</para>
</section>
<section id="channel-configuration-executorchannel">
<title>ExecutorChannel</title>
<para>
To create an <classname>ExecutorChannel</classname>, add the 'task-executor' attribute. Its value
can reference any <interfacename>TaskExecutor</interfacename> within the context. For example,
this enables configuration of a thread-pool for dispatching messages to subscribed handlers.
As mentioned above, this does break the "single-threaded" execution context between sender
and receiver so that any active transaction context will not be shared by the invocation
of the handler (i.e. the handler may throw an Exception, but the send invocation has already
returned successfully).
<programlisting language="xml">&lt;channel id="executorChannel" task-executor="someExecutor"/&gt;</programlisting>
</para>
</section>
<section id="channel-configuration-prioritychannel">
<title>PriorityChannel Configuration</title>
<para>
To create a <classname>PriorityChannel</classname>, use the "priority-queue" sub-element:
<programlisting language="xml"><![CDATA[<channel id="exampleChannel">
<programlisting language="xml"><![CDATA[<channel id="priorityChannel">
<priority-queue capacity="20"/>
</channel>]]></programlisting>
By default, the channel will consult the <classname>MessagePriority</classname> header of the
@@ -338,7 +423,7 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
provided instead. Also, note that the <classname>PriorityChannel</classname> (like the other types)
does support the "datatype" attribute. As with the QueueChannel, it also supports a "capacity" attribute.
The following example demonstrates all of these:
<programlisting language="xml"><![CDATA[<channel id="exampleChannel" datatype="example.Widget">
<programlisting language="xml"><![CDATA[<channel id="priorityChannel" datatype="example.Widget">
<priority-queue comparator="widgetComparator"
capacity="10"/>
</channel>
@@ -349,8 +434,10 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
<title>RendezvousChannel Configuration</title>
<para>
A <classname>RendezvousChannel</classname> is created when the queue sub-element is
a &lt;rendezvous-queue&gt;. It does not provide any additional configuration options.
<programlisting language="xml"><![CDATA[<channel id="exampleChannel"/>
a &lt;rendezvous-queue&gt;. It does not provide any additional configuration options to
those described above, and its queue does not accept any capacity value since it is a
0-capacity direct handoff queue.
<programlisting language="xml"><![CDATA[<channel id="rendezvousChannel"/>
<rendezvous-queue/>
</channel>
]]></programlisting>
@@ -360,7 +447,7 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
<title>ThreadLocalChannel Configuration</title>
<para>
The <classname>ThreadLocalChannel</classname> does not provide any additional configuration options.
<programlisting language="xml"><![CDATA[<thread-local-channel id="exampleChannel"/>]]></programlisting>
<programlisting language="xml"><![CDATA[<thread-local-channel id="threadLocalChannel"/>]]></programlisting>
</para>
</section>
<para>
@@ -376,5 +463,17 @@ public Message<?> receive(final PollableChannel<?> channel) { ... }]]></programl
In general, it is a good idea to define the interceptor implementations in a separate location since they
usually provide common behavior that can be reused across multiple channels.
</para>
<note>
<para>
If namespace support is enabled, there are also two special channels defined within the context by default:
<code>errorChannel</code> and <code>nullChannel</code>. The 'nullChannel' acts like <code>/dev/null</code>,
simply logging any Message sent to it at DEBUG level and returning immediately. Any time you face channel
resolution errors for a reply that you don't care about, you can set the affected component's 'output-channel'
to reference 'nullChannel' (the name 'nullChannel' is reserved within the context). The 'errorChannel' is
used internally for sending error messages, and it can be overridden with a custom configuration. It is
discussed in greater detail in <xref linkend="namespace-errorhandler"/>.
</para>
</note>
</section>
</chapter>