diff --git a/spring-integration-reference/src/channel.xml b/spring-integration-reference/src/channel.xml index c4ebd47e4d..4f105f6de4 100644 --- a/spring-integration-reference/src/channel.xml +++ b/spring-integration-reference/src/channel.xml @@ -22,11 +22,14 @@ 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 { + +
+ PollableChannel + + 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(); @@ -37,21 +40,26 @@ 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 { + Similar to the send methods, when receiving a message, the return value will be null in the + case of a timeout or interrupt. + +
+ +
+ SubscribableChannel + + 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); } - + +
diff --git a/spring-integration-reference/src/core-api.xml b/spring-integration-reference/src/core-api.xml index 63de737885..3079fa1d79 100644 --- a/spring-integration-reference/src/core-api.xml +++ b/spring-integration-reference/src/core-api.xml @@ -1,371 +1,4 @@ - - - - The Core API -
- Message - - 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. Here is the definition of the - Message interface: - public interface Message<T> { - - T getPayload(); - - MessageHeaders getHeaders(); - -} - And the following 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/outbound adapter implementations will also provide and/or expect certain headers, and additional - user-defined headers can also be configured. - - - 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 latter accepts any - Throwable object as its payload. - - - 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 section). It is defined as an enum with five possible values: - public enum MessagePriority { - HIGHEST, - HIGH, - NORMAL, - LOW, - LOWEST -} - - - 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. - -
- -
- MessageChannel - - While the Message plays the crucial role of encapsulating data, it is the - MessageChannel that decouples message producers from message consumers. - 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); - -} - - - 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. - -
-
- -
- ChannelInterceptor - - 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; - } -}]]> - -
MessageHandler @@ -498,50 +131,6 @@ channel.addInterceptor(interceptor);
-
- MessageExchangeTemplate - - Whereas the MessageHandler interface provides the foundation for many of the - components that enable non-invasive invocation of your application code from the messaging - system, sometimes it is necessary to invoke the messaging system from your application - code. Spring Integration provides a MessageExchangeTemplate that supports a - variety of message-exchanges, including request/reply scenarios. For example, it is possible to send a request - and wait for a reply. - MessageExchangeTemplate template = new MessageExchangeTemplate(); -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 MessageTarget target) { ... } - -public Message sendAndReceive(final Message request, final MessageTarget target) { .. } - -public Message receive(final PollableSource source) { ... } - -public boolean receiveAndForward(final PollableSource source, final MessageTarget target) { ... }]]> - - - Additionally, a 'transactionManager' can be configured on a MessageExchangeTemplate as well as the various - transaction attributes: - template.setTransactionManager(transactionManager); -template.setPropagationBehaviorName(propagationBehavior); -template.setIsolationLevelName(isolationLevel); -template.setTransactionTimeout(transactionTimeout); -template.setTransactionReadOnly(readOnly); -template.setReceiveTimeout(receiveTimeout); -template.setSendTimeout(sendTimeout); - Finally, there is a also an asynchronous version called AsyncMessageExchangeTemplate - whose constructor accepts a TaskExecutor, and whose Message-returning methods - return an AsyncMessage. That is essentially a wrapper for any Message that also - implements Future<Message<T>>: - AsyncMessageExchangeTemplate template = new AsyncMessageExchangeTemplate(taskExecutor); -Message reply = template.sendAndReceive(new StringMessage("test"), someChannel); -// do some work in the meantime -reply.getPayload(); // blocks if still waiting for actual reply - - -
-
MessagingGateway diff --git a/spring-integration-reference/src/endpoint.xml b/spring-integration-reference/src/endpoint.xml new file mode 100644 index 0000000000..bd34c4758d --- /dev/null +++ b/spring-integration-reference/src/endpoint.xml @@ -0,0 +1,119 @@ + + + + Message Endpoints + + As mentioned in the overview, Message Endpoints are responsible for connecting the various messaging components to + channels. Over the next several chapters, you will see a number of different components that consume Messages. Some + of these are also capable of sending reply Messages. Sending Messages is quite straightforward. As shown above in + , it's easy to send a Message to a Message Channel. However, + receiving is a bit more complicated. The main reason is that there are two types of consumers: + Polling Consumers and + Event-Driven Consumers. + + + Of the two, Event-Driven Consumers are much simpler. Without any need to manage and schedule a separate poller + thread, they are essentially just listeners with a callback method. When connecting to one of Spring Integration's + subscribable Message Channels, this simple option works great. However, when connecting to a buffering, pollable + Message Channel, some component has to schedule and manage the polling thread(s). Spring Integration provides + two different endpoint implementations to accommodate these two types of consumers. Therefore, the consumers + themselves can simply implement the callback interface. When polling is required, the endpoint acts as a + "container" for the consumer instance. The benefit is similar to that of using a container for hosting + Message-Driven Beans, but since these consumers are simply Spring-managed Objects running within an + ApplicationContext, it more closely resembles Spring's own MessageListener containers. + + +
+ Message Consumer + + Spring Integration's MessageConsumer interface is defined as follows: + public interface MessageConsumer { + + void onMessage(Message<?> message); + +} + Despite its simplicity, this provides the foundation for most of the components that will be covered in the + following chapters (Routers, Transformers, Splitters, Aggregators, Service Activators, etc). Those components + each perform very different functionality with the Messages they receive, but the requirements for actually + receiving a Message are the same, and the choice between polling and event-driven behavior is also the same. + Spring Integration provides two endpoint implementations that "host" these callback-based consumers and allow + them to be connected to Message Channels. + +
+ +
+ Event-Driven Consumer + + Because it is the simpler of the two, we will cover the Event-Driven Consumer endpoint first. You may recall that + the SubscribableChannel interface provides a subscribe() + method and that the method accepts a MessageConsumer parameter (as shown in + ): + +subscribableChannel.subscribe(messageConsumer); + + Since a consumer that is subscribed to a channel does not have to actively poll that channel, this is an + Event-Driven Consumer, and the corresponding endpoint "container" class provided by Spring Integration accepts a + MessageConsumer and a SubscribableChannel: + MessageConsumer consumer = new ExampleConsumer(); + +SubscribableChannel channel = (SubscribableChannel) context.getBean("exampleSubscribableChannel"); + +SubscribingConsumerEndpoint endpoint = new SubscribingConsumerEndpoint(consumer, channel); + +
+ +
+ Polling Consumer + + Spring Integration also provides a PollingConsumerEndpoint, and it can be instantiated in + the same way except that the channel must implement PollableChannel: + MessageConsumer consumer = new ExampleConsumer(); + +PollableChannel channel = (PollableChannel) context.getBean("examplePollableChannel"); + +PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(consumer, channel); + + + There are many other configuration options for the polling endpoint. For example, the trigger can be provided: + +PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(consumer, channel); + +endpoint.setTrigger(new IntervalTrigger(30, TimeUnit.SECONDS)); + Likewise, other polling-related configuration properties may be specified: + +PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(consumer, channel); + +endpoint.setMaxMessagesPerPoll(10); + +endpoint.setReceiveTimeout(5000); + A polling consumer may even delegate to a Spring TaskExecutor and + participate in Spring-managed transactions. The following example shows the configuration of both: + +PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(consumer, channel); + +TaskExecutor taskExecutor = (TaskExecutor) context.getBean("exampleExecutor"); +endpoint.setTaskExecutor(taskExecutor); + +PlatformTransactionManager txManager = (PlatformTransationManager) context.getBean("exampleTxManager"); +endpoint.setTransactionManager(txManager); + The examples above show dependency lookups, but keep in mind that these endpoints will most often be configured + as Spring bean definitions. In fact, Spring Integration also provides a + FactoryBean that creates the appropriate endpoint type based on the type of + channel, and there is full XML namespace support to even further hide those details. The namespace-based + configuration will be featured as each component type is introduced. + + Interestingly, many of the MessageConsumer implementations are also capable of + generating reply Messages. As mentioned above, sending Messages is trivial when compared to the Message + reception. Nevertheless, when and how many reply Messages are sent + depends on the consumer type. For example, an Aggregator waits for a number of Messages to + arrive and is often a downstream consumer for a Splitter which may generate multiple + replies for each Message it consumes. When using the namespace configuration, you do not strictly need to know + all of the details, but it still might be worth knowing that several of these components share a common base + class, the AbstractReplyProducingMessageConsumer, and it provides a + setOutputChannel(..) method. + + +
+ +
\ 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 f174ad1446..2384dad709 100644 --- a/spring-integration-reference/src/spring-integration-reference.xml +++ b/spring-integration-reference/src/spring-integration-reference.xml @@ -43,8 +43,9 @@ - + +