diff --git a/spring-integration-reference/src/channel.xml b/spring-integration-reference/src/channel.xml new file mode 100644 index 0000000000..c4ebd47e4d --- /dev/null +++ b/spring-integration-reference/src/channel.xml @@ -0,0 +1,242 @@ + + + + Message Channels + + While the Message plays the crucial role of encapsulating data, it is the + MessageChannel that decouples message producers from message consumers. + + +
+ The MessageChannel Interface + + Spring Integration's top-level MessageChannel interface is defined as follows. + + When sending a message, the return value will be true if the message is sent successfully. + If the send call times out or is interrupted, then it will return false. + + + Since Message Channels may or may not buffer Messages (as discussed in the overview), there are two + sub-interfaces defining the buffering (pollable) and non-buffering (subscribable) channel behavior. Here is the + definition of PollableChannel. + public interface PollableChannel extends MessageChannel { + + Message<?> receive(); + + Message<?> receive(long timeout); + + List<Message<?>> clear(); + + List<Message<?>> purge(MessageSelector selector); + +} + Similar to the send methods, when receiving a message, the return value will be null in the + case of a timeout or interrupt. + + + The SubscribableChannel base interface is implemented by channels that send + Messages directly to their subscribed consumers. Therefore, they do not provide receive methods for polling, but + instead define methods for handling those subscribers: + public interface SubscribableChannel extends MessageChannel { + + boolean subscribe(MessageConsumer consumer); + + boolean unsubscribe(MessageConsumer consumer); + +} + +
+ +
+ Message Channel Implementations + + Spring Integration provides several different Message Channel implementations. Each is briefly described in the + sections below. + +
+ PublishSubscribeChannel + + The PublishSubscribeChannel implementation broadcasts any Message + sent to it to all of its subscribed consumers. This is most often used for sending + Event Messages whose primary role is notification as opposed to + Document Messages which are generally intended to be processed by + a single consumer. Note that the PublishSubscribeChannel is + intended for sending only. Since it broadcasts to its subscribers directly when its + send(Message) method is invoked, consumers cannot poll for + Messages (it does not implement PollableChannel and + therefore has no receive() method). Instead, any subscriber + must be a MessageConsumer itself, and the subscriber's + send(Message) method will be invoked in turn. + +
+
+ QueueChannel + + The QueueChannel implementation wraps a queue. Unlike, the + PublishSubscribeChannel, the QueueChannel has point-to-point + semantics. In other words, even if the channel has multiple consumers, only one of them should receive any + Message sent to that channel. It provides a default no-argument constructor (providing an essentially unbounded + capacity of Integer.MAX_VALUE) as well as a constructor that accepts the queue capacity: + public QueueChannel(int capacity) + A channel that has not reached its capacity limit will store messages in its internal queue, and the + send() 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. + Note however, that calling the no-arg versions of send() and + receive() will block indefinitely. + +
+
+ PriorityChannel + + Whereas the QueueChannel enforces first-in/first-out (FIFO) ordering, the + PriorityChannel 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 + 'priority' header within each message. However, for custom priority determination + logic, a comparator of type Comparator<Message<?>> can be provided to the + PriorityChannel's constructor. + +
+
+ RendezvousChannel + + The RendezvousChannel enables a "direct-handoff" scenario where a sender will block + until another party invokes the channel's receive() method or vice-versa. Internally, + this implementation is quite similar to the QueueChannel except that it uses a + SynchronousQueue (a zero-capacity implementation of + BlockingQueue). 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 + QueueChannel, the message would have been stored to the internal queue and potentially + never received. + + + The RendezvousChannel is also useful for implementing request-reply + operations. The sender can create a temporary, anonymous instance of RendezvousChannel + 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. + +
+
+ DirectChannel + + The DirectChannel has point-to-point semantics, but otherwise is more similar to the + PublishSubscribeChannel than any of the queue-based channel implementations described + above. It implements the SubscribableChannel interface instead of the + PollableChannel interface, so it dispatches Messages directly to a subscriber. + As a point-to-point channel, however, it differs from the PublishSubscribeChannel in + that it will only send each Message to a single subscribed + MessageConsumer. Its primary purpose is to enable a single thread to perform the + operations on "both sides" of the channel. For example, if a consumer is subscribed to a + DirectChannel, then sending a Message to that channel will trigger invocation of that + consumer's onMessage(Message) method directly in the sender's + thread. 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 consumer's invocation (e.g. updating a database record) can play a role in determining the ultimate result + of that transaction (commit or rollback). + + Since the DirectChannel 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 PollableChannels. Likewise, if a channel needs to broadcast + messages, it should not be a DirectChannel but rather a + PublishSubscribeChannel. Below you will see how these can be configured. + + +
+
+ ThreadLocalChannel + + The final channel implementation type is ThreadLocalChannel. This channel also delegates + to a queue internally, but the queue is bound to the current thread. That way the thread that sends to the + channel will later be able to receive those same Messages, but no other thread would be able to access them. + While probably the least common type of channel, this is useful for situations where + DirectChannels are being used to enforce a single thread of operation but any reply + Messages should be sent to a "terminal" channel. If that terminal channel is a + ThreadLocalChannel, the original sending thread can collect its replies from it. + +
+
+ +
+ Channel Interceptors + + One of the advantages of a messaging architecture is the ability to provide common behavior and capture + meaningful information about the messages passing through the system in a non-invasive way. Since the + Messages are being sent to and received from + MessageChannels, those channels provide an opportunity for intercepting + the send and receive operations. The ChannelInterceptor strategy interface + provides methods for each of those operations: + preSend(Message message, MessageChannel channel); + + void postSend(Message message, MessageChannel channel, boolean sent); + + boolean preReceive(MessageChannel channel); + + Message postReceive(Message message, MessageChannel channel); +}]]> + After implementing the interface, registering the interceptor with a channel is just a matter of calling: + channel.addInterceptor(someChannelInterceptor); + 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 + preReceive method can return 'false' to prevent the receive + operation from proceeding. + + + Because it is rarely necessary to implement all of the interceptor methods, a + ChannelInterceptorAdapter class is also available for sub-classing. It provides no-op + methods (the void method is empty, the Message returning methods + return the Message parameter as-is, and the boolean method returns true). + Therefore, it is often easiest to extend that class and just implement the method(s) that you need as in the + following example. + preSend(Message message, MessageChannel channel) { + sendCount.incrementAndGet(); + return message; + } +}]]> + +
+ +
+ MessageChannelTemplate + + As you will see when the endpoints and their various configuration options are introduced, Spring Integration + provides a foundation for messaging components that enables non-invasive invocation of your application code + from the messaging system. However, sometimes it is necessary to invoke the messaging system + from your application code. For convenience when implementing such use-cases, Spring + Integration provides a MessageChannelTemplate that supports a variety of operations across + the Message Channels, including request/reply scenarios. For example, it is possible to send a request + and wait for a reply. + MessageChannelTemplate template = new MessageChannelTemplate(); + +Message reply = template.sendAndReceive(new StringMessage("test"), someChannel); + In that example, a temporary anonymous channel would be created internally by the template. The + 'sendTimeout' and 'receiveTimeout' properties may also be set on the template, and other exchange + types are also supported. + message, final MessageChannel channel) { ... } + +public Message sendAndReceive(final Message request, final MessageChannel channel) { .. } + +public Message receive(final PollableChannel channel) { ... }]]> + +
+
\ No newline at end of file diff --git a/spring-integration-reference/src/message.xml b/spring-integration-reference/src/message.xml new file mode 100644 index 0000000000..c2427f0911 --- /dev/null +++ b/spring-integration-reference/src/message.xml @@ -0,0 +1,218 @@ + + + + Message Construction + + The Spring Integration Message is a generic container for data. Any object can + be provided as the payload, and each Message also includes headers containing + user-extensible properties as key-value pairs. + + +
+ The Message Interface + Here is the definition of the Message interface: + public interface Message<T> { + + T getPayload(); + + MessageHeaders getHeaders(); + +} + + + The Message is obviously a very important part of the API. By encapsulating the + data in a generic wrapper, the messaging system can pass it around without any knowledge of the data's type. As + an application evolves to support new types, or when the types themselves are modified and/or extended, the + messaging system will not be affected by such changes. On the other hand, when some component in the messaging + system does require access to information about the Message, + such metadata can typically be stored to and retrieved from the metadata in the Message Headers. + +
+ +
+ Message Headers + + Just as Spring Integration allows any Object to be used as the payload of a Message, it also supports any Object + types as header values. In fact, the MessageHeaders class implements the + java.util.Map interface: + public final class MessageHeaders implements Map<String, Object>, Serializable { + ... +} + + Even though the MessageHeaders implements Map, it is effectively a read-only implementation. Any attempt to + put a value in the Map will result in an UnsupportedOperationException. + The same applies for remove and clear. Since Messages may be passed to + multiple consumers, the structure of the Map cannot be modified. Likewise, the Message's payload Object can not + be set after the initial creation. However, the mutability of the header values themselves + (or the payload Object) is intentionally left as a decision for the framework user. + + + + As an implementation of Map, the headers can obviously be retrieved by calling get(..) + with the name of the header. Alternatively, you can provide the expected Class as an + additional parameter. Even better, when retrieving one of the pre-defined values, convenient getters are + available. Here is an example of each of these three options: + Object someValue = message.getHeaders().get("someKey"); + + CustomerId customerId = message.getHeaders().get("customerId", CustomerId.class); + + Long timestamp = message.getHeaders().getTimestamp(); + + + + The following Message headers are pre-defined: + + Pre-defined Message Headers + + + + + Header Name + Header Type + + + + + ID + java.util.UUID + + + TIMESTAMP + java.lang.Long + + + EXPIRATION_DATE + java.lang.Long + + + CORRELATION_ID + java.lang.Object + + + REPLY_CHANNEL + java.lang.Object (can be a String or MessageChannel) + + + SEQUENCE_NUMBER + java.lang.Integer + + + SEQUENCE_SIZE + java.lang.Integer + + + PRIORITY + MessagePriority (an enum) + + + +
+
+ + Many inbound and outbound adapter implementations will also provide and/or expect certain headers, and additional + user-defined headers can also be configured. + +
+ +
+ Message Implementations + + The base implementation of the Message interface is + GenericMessage<T>, and it provides two constructors: + new GenericMessage<T>(T payload); + +new GenericMessage<T>(T payload, Map<String, Object> headers) + When a Message is created, a random unique id will be generated. The constructor that accepts a Map of headers + will copy the provided headers to the newly created Message. + + + There are also two convenient subclasses available: StringMessage and + ErrorMessage. The former accepts a String as its payload: + StringMessage message = new StringMessage("hello world"); + +String s = message.getPayload(); + And, the latter accepts any Throwable object as its payload: + ErrorMessage message = new ErrorMessage(someThrowable); + +Throwable t = message.getPayload(); + Notice that these implementations take advantage of the fact that the GenericMessage + base class is parameterized. Therefore, as shown in both examples, no casting is necessary when retrieving + the Message payload Object. + +
+ +
+ The MessageBuilder Helper Class + + You may notice that the Message interface defines retrieval methods for its payload and headers but no setters. + The reason for this is that a Message cannot be modified after its initial creation. Therefore, when a Message + instance is sent to multiple consumers (e.g. through a Publish Subscribe Channel), if one of those consumers + needs to send a reply with a different payload type, it will need to create a new Message. As a result, the + other consumers are not affected by those changes. Keep in mind, that multiple consumers may access the same + payload instance or header value, and whether such an instance is itself immutable is a decision left to the + developer. In other words, the contract for Messages is similar to that of an + unmodifiable Collection, and the MessageHeaders' map further exemplifies that; even though + the MessageHeaders class implements java.util.Map, any attempt to invoke a + put operation (or 'remove' or 'clear') on the MessageHeaders will result in an + UnsupportedOperationException. + + + Rather than requiring the creation and population of a Map to pass into the GenericMessage constructor, Spring + Integration does provide a far more convenient way to construct Messages: MessageBuilder. + The MessageBuilder provides two factory methods for creating Messages from either an existing Message or with a + payload Object. When building from an existing Message, the headers and payload of that + Message will be copied to the new Message: + Message<String> message1 = MessageBuilder.withPayload("test") + .setHeader("foo", "bar") + .build(); + +Message<String> message2 = MessageBuilder.fromMessage(message1).build(); + +assertEquals("test", message2.getPayload()); +assertEquals("bar", message2.getHeaders().get("foo")); + + + If you need to create a Message with a new payload but still want to copy the + headers from an existing Message, you can use one of the 'copy' methods. + Message<String> message3 = MessageBuilder.fromPayload("test3") + .copyHeaders(message1.getHeaders()) + .build(); + +Message<String> message4 = MessageBuilder.fromPayload("test4") + .setHeader("foo", 123) + .copyHeadersIfAbsent(message1.getHeaders()) + .build(); + +assertEquals("bar", message3.getHeaders().get("foo")); +assertEquals(123, message4.getHeaders().get("foo")); + Notice that the copyHeadersIfAbsent does not overwrite existing values. Also, in the + second example above, you can see how to set any user-defined header with setHeader. + Finally, there are set methods available for the predefined headers as well as a non-destructive method for + setting any header (MessageHeaders also defines constants for the pre-defined header names). + Message<Integer> importantMessage = MessageBuilder.fromPayload(99) + .setPriority(MessagePriority.HIGHEST) + .build(); + +assertEquals(MessagePriority.HIGHEST, importantMessage.getHeaders().getPriority()); + +Message<Integer> anotherMessage = MessageBuilder.fromMessage(importantMessage) + .setHeaderIfAbsent(MessageHeaders.PRIORITY, MessagePriority.LOW) + .build(); + +assertEquals(MessagePriority.HIGHEST, anotherMessage.getHeaders().getPriority()); + + + + The MessagePriority is only considered when using a PriorityChannel + (as described in the next chapter). It is defined as an enum with five possible values: + public enum MessagePriority { + HIGHEST, + HIGH, + NORMAL, + LOW, + LOWEST +} + +
+ +
\ No newline at end of file diff --git a/spring-integration-reference/src/spring-integration-reference.xml b/spring-integration-reference/src/spring-integration-reference.xml index 750ee8fcd3..f174ad1446 100644 --- a/spring-integration-reference/src/spring-integration-reference.xml +++ b/spring-integration-reference/src/spring-integration-reference.xml @@ -41,7 +41,8 @@ - + +