diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/InboundChannelAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/InboundChannelAdapter.java
new file mode 100644
index 0000000000..0f352539dd
--- /dev/null
+++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/InboundChannelAdapter.java
@@ -0,0 +1,66 @@
+/*
+ * Copyright 2014 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.integration.annotation;
+
+import java.lang.annotation.Documented;
+import java.lang.annotation.ElementType;
+import java.lang.annotation.Inherited;
+import java.lang.annotation.Retention;
+import java.lang.annotation.RetentionPolicy;
+import java.lang.annotation.Target;
+
+/**
+ * Indicates that a method is capable of producing a {@link org.springframework.messaging.Message}
+ * or {@link org.springframework.messaging.Message} {@code payload}.
+ *
+ * A method annotated with {@code @InboundChannelAdapter} can't accept any parameters.
+ *
+ * Return values from the annotated method may be of any type. If the return
+ * value is not a {@link org.springframework.messaging.Message}, a {@link org.springframework.messaging.Message}
+ * will be created with that object as its {@code payload}.
+ *
+ * The result {@link org.springframework.messaging.Message} will be sent to the provided {@link #value()}.
+ *
+ * {@code @InboundChannelAdapter} is an analogue of {@code }. With that
+ * the {@link org.springframework.integration.scheduling.PollerMetadata} is required to to initiate
+ * the method invocation. Or {@link #poller()} should be provided, or the
+ * {@link org.springframework.integration.scheduling.PollerMetadata#DEFAULT_POLLER} bean has to be configured
+ * in the application context.
+ *
+ *
+ * @author Artem Bilan
+ * @since 4.0
+ */
+@Target(ElementType.METHOD)
+@Retention(RetentionPolicy.RUNTIME)
+@Inherited
+@Documented
+public @interface InboundChannelAdapter {
+
+ /**
+ * @return the 'channel' bean name to send the {@link org.springframework.messaging.Message}.
+ */
+ String value();
+
+ /**
+ * @return the {@link org.springframework.integration.annotation.Poller} options for a polled endpoint
+ * ({@link org.springframework.integration.scheduling.PollerMetadata}).
+ * This attribute is an {@code array} just to allow an empty default (no poller).
+ * Only one {@link org.springframework.integration.annotation.Poller} element is allowed.
+ */
+ Poller[] poller() default {};
+}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java
index f4b45e0df4..8279aefc53 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java
@@ -37,6 +37,7 @@ import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.config.IntegrationConfigUtils;
import org.springframework.integration.context.Orderable;
import org.springframework.integration.endpoint.AbstractEndpoint;
+import org.springframework.integration.endpoint.AbstractPollingEndpoint;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.endpoint.PollingConsumer;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
@@ -159,74 +160,8 @@ public abstract class AbstractMethodAnnotationPostProcessor annotationType) {
+ protected void configurePollingEndpoint(AbstractPollingEndpoint pollingEndpoint, T annotation) {
+ PollerMetadata pollerMetadata = null;
+ Poller[] pollers = (Poller[]) AnnotationUtils.getValue(annotation, "poller");
+ if (!ObjectUtils.isEmpty(pollers)) {
+ Assert.state(pollers.length == 1, "The 'poller' for an Annotation-based endpoint can have only one '@Poller'.");
+ Poller poller = pollers[0];
+
+ String ref = poller.value();
+ String triggerRef = poller.trigger();
+ String executorRef = poller.taskExecutor();
+ String fixedDelayValue = this.environment.resolvePlaceholders(poller.fixedDelay());
+ String fixedRateValue = this.environment.resolvePlaceholders(poller.fixedRate());
+ String maxMessagesPerPollValue = this.environment.resolvePlaceholders(poller.maxMessagesPerPoll());
+ String cron = this.environment.resolvePlaceholders(poller.cron());
+
+ if (StringUtils.hasText(ref)) {
+ Assert.state(!StringUtils.hasText(triggerRef) && !StringUtils.hasText(executorRef) && !StringUtils.hasText(cron)
+ && !StringUtils.hasText(fixedDelayValue) && !StringUtils.hasText(fixedRateValue)
+ && !StringUtils.hasText(maxMessagesPerPollValue),
+ "The '@Poller' 'ref' attribute is mutually exclusive with other attributes.");
+ pollerMetadata = this.beanFactory.getBean(ref, PollerMetadata.class);
+ }
+ else {
+ pollerMetadata = new PollerMetadata();
+ if (StringUtils.hasText(maxMessagesPerPollValue)) {
+ pollerMetadata.setMaxMessagesPerPoll(Long.parseLong(maxMessagesPerPollValue));
+ }
+ if (StringUtils.hasText(executorRef)) {
+ pollerMetadata.setTaskExecutor(this.beanFactory.getBean(executorRef, TaskExecutor.class));
+ }
+ Trigger trigger = null;
+ if (StringUtils.hasText(triggerRef)) {
+ Assert.state(!StringUtils.hasText(cron) && !StringUtils.hasText(fixedDelayValue) && !StringUtils.hasText(fixedRateValue),
+ "The '@Poller' 'trigger' attribute is mutually exclusive with other attributes.");
+ trigger = this.beanFactory.getBean(triggerRef, Trigger.class);
+ }
+ else if (StringUtils.hasText(cron)) {
+ Assert.state(!StringUtils.hasText(fixedDelayValue) && !StringUtils.hasText(fixedRateValue),
+ "The '@Poller' 'cron' attribute is mutually exclusive with other attributes.");
+ trigger = new CronTrigger(cron);
+ }
+ else if (StringUtils.hasText(fixedDelayValue)) {
+ Assert.state(!StringUtils.hasText(fixedRateValue),
+ "The '@Poller' 'fixedDelay' attribute is mutually exclusive with other attributes.");
+ trigger = new PeriodicTrigger(Long.parseLong(fixedDelayValue));
+ }
+ else if (StringUtils.hasText(fixedRateValue)) {
+ trigger = new PeriodicTrigger(Long.parseLong(fixedRateValue));
+ ((PeriodicTrigger) trigger).setFixedRate(true);
+ }
+ //'Trigger' can be null. 'PollingConsumer' does fallback to the 'new PeriodicTrigger(10)'.
+ pollerMetadata.setTrigger(trigger);
+ }
+
+
+ }
+ else {
+ pollerMetadata = PollerMetadata.getDefaultPollerMetadata(this.beanFactory);
+ Assert.notNull(pollerMetadata, "No poller has been defined for Annotation-based endpoint, " +
+ "and no default poller is available within the context.");
+ }
+ pollingEndpoint.setTaskExecutor(pollerMetadata.getTaskExecutor());
+ pollingEndpoint.setTrigger(pollerMetadata.getTrigger());
+ pollingEndpoint.setAdviceChain(pollerMetadata.getAdviceChain());
+ pollingEndpoint.setMaxMessagesPerPoll(pollerMetadata.getMaxMessagesPerPoll());
+ pollingEndpoint.setErrorHandler(pollerMetadata.getErrorHandler());
+ if (pollingEndpoint instanceof PollingConsumer) {
+ ((PollingConsumer) pollingEndpoint).setReceiveTimeout(pollerMetadata.getReceiveTimeout());
+ }
+ pollingEndpoint.setTransactionSynchronizationFactory(pollerMetadata.getTransactionSynchronizationFactory());
+ }
+
+ protected String generateHandlerBeanName(String originalBeanName, Method method, Class extends Annotation> annotationType) {
String baseName = originalBeanName + "." + method.getName() + "." + ClassUtils.getShortNameAsProperty(annotationType);
String name = baseName;
int count = 1;
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java
new file mode 100644
index 0000000000..96837ae8b5
--- /dev/null
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java
@@ -0,0 +1,90 @@
+/*
+ * Copyright 2014 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.integration.config.annotation;
+
+import java.lang.annotation.Annotation;
+import java.lang.reflect.Method;
+
+import org.springframework.beans.factory.ListableBeanFactory;
+import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
+import org.springframework.core.annotation.AnnotationUtils;
+import org.springframework.core.env.Environment;
+import org.springframework.integration.annotation.InboundChannelAdapter;
+import org.springframework.integration.config.IntegrationConfigUtils;
+import org.springframework.integration.endpoint.MethodInvokingMessageSource;
+import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
+import org.springframework.messaging.MessageChannel;
+import org.springframework.messaging.MessageHandler;
+import org.springframework.util.Assert;
+
+/**
+ * Post-processor for Methods annotated with {@link InboundChannelAdapter @InboundChannelAdapter}.
+ *
+ * @author Artem Bilan
+ * @since 4.0
+ */
+public class InboundChannelAdapterAnnotationPostProcessor extends
+ AbstractMethodAnnotationPostProcessor {
+
+ public InboundChannelAdapterAnnotationPostProcessor(ListableBeanFactory beanFactory, Environment environment) {
+ super(beanFactory, environment);
+ }
+
+ @Override
+ public Object postProcess(Object bean, String beanName, Method method, InboundChannelAdapter annotation) {
+ Assert.isTrue(!Void.class.isAssignableFrom(method.getReturnType()), "The method '" + method
+ + "' for 'SourcePollingChannelAdapter' must not have 'void' return type.");
+ Assert.isTrue(method.getParameterTypes().length == 0, "The method '" + method
+ + "' for 'SourcePollingChannelAdapter' must not have any parameters.");
+
+ String channelName = (String) AnnotationUtils.getValue(annotation);
+ Assert.hasText(channelName, "The channel ('value' attribute of @InboundChannelAdapter) can't be empty.");
+
+ MessageChannel channel = this.channelResolver.resolveDestination(channelName);
+
+ MethodInvokingMessageSource messageSource = new MethodInvokingMessageSource();
+ messageSource.setObject(bean);
+ messageSource.setMethod(method);
+ if (beanFactory instanceof ConfigurableListableBeanFactory) {
+ String handlerBeanName = this.generateHandlerBeanName(beanName, method, annotation.annotationType());
+ ConfigurableListableBeanFactory listableBeanFactory = (ConfigurableListableBeanFactory) beanFactory;
+ listableBeanFactory.registerSingleton(handlerBeanName, messageSource);
+ messageSource = (MethodInvokingMessageSource) listableBeanFactory
+ .initializeBean(messageSource, handlerBeanName);
+ }
+
+ SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter();
+ adapter.setOutputChannel(channel);
+ adapter.setSource(messageSource);
+ this.configurePollingEndpoint(adapter, annotation);
+
+ return adapter;
+ }
+
+ @Override
+ protected String generateHandlerBeanName(String originalBeanName, Method method,
+ Class extends Annotation> annotationType) {
+ return super.generateHandlerBeanName(originalBeanName, method, annotationType)
+ .replaceFirst(IntegrationConfigUtils.HANDLER_ALIAS_SUFFIX + "$", ".source");
+ }
+
+ @Override
+ protected MessageHandler createHandler(Object bean, Method method, InboundChannelAdapter annotation) {
+ throw new UnsupportedOperationException();
+ }
+
+}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java
index 03559384e1..dec6f2d04c 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java
@@ -34,7 +34,6 @@ import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.beans.factory.BeanInitializationException;
-import org.springframework.beans.factory.BeanNameAware;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.config.BeanPostProcessor;
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
@@ -42,11 +41,11 @@ import org.springframework.context.ApplicationEvent;
import org.springframework.context.ApplicationListener;
import org.springframework.context.EnvironmentAware;
import org.springframework.context.Lifecycle;
-import org.springframework.context.SmartLifecycle;
import org.springframework.core.annotation.AnnotationUtils;
import org.springframework.core.env.Environment;
import org.springframework.integration.annotation.Aggregator;
import org.springframework.integration.annotation.Filter;
+import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.Router;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.annotation.Splitter;
@@ -104,6 +103,7 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean
postProcessors.put(ServiceActivator.class, new ServiceActivatorAnnotationPostProcessor(this.beanFactory, this.environment));
postProcessors.put(Splitter.class, new SplitterAnnotationPostProcessor(this.beanFactory, this.environment));
postProcessors.put(Aggregator.class, new AggregatorAnnotationPostProcessor(this.beanFactory, this.environment));
+ postProcessors.put(InboundChannelAdapter.class, new InboundChannelAdapterAnnotationPostProcessor(this.beanFactory, this.environment));
}
public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException {
@@ -126,27 +126,20 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean
if (postProcessor != null && shouldCreateEndpoint(annotation)) {
Object result = postProcessor.postProcess(bean, beanName, method, annotation);
if (result != null && result instanceof AbstractEndpoint) {
+ AbstractEndpoint endpoint = (AbstractEndpoint) result;
String endpointBeanName = generateBeanName(beanName, method, annotation.annotationType());
- if (result instanceof BeanNameAware) {
- ((BeanNameAware) result).setBeanName(endpointBeanName);
+ endpoint.setBeanName(endpointBeanName);
+ beanFactory.registerSingleton(endpointBeanName, endpoint);
+ endpoint.setBeanFactory(beanFactory);
+ try {
+ endpoint.afterPropertiesSet();
}
- beanFactory.registerSingleton(endpointBeanName, result);
- if (result instanceof BeanFactoryAware) {
- ((BeanFactoryAware) result).setBeanFactory(beanFactory);
+ catch (Exception e) {
+ throw new BeanInitializationException("failed to initialize annotated component", e);
}
- if (result instanceof InitializingBean) {
- try {
- ((InitializingBean) result).afterPropertiesSet();
- }
- catch (Exception e) {
- throw new BeanInitializationException("failed to initialize annotated component", e);
- }
- }
- if (result instanceof Lifecycle) {
- lifecycles.add((Lifecycle) result);
- if (result instanceof SmartLifecycle && ((SmartLifecycle) result).isAutoStartup()) {
- ((SmartLifecycle) result).start();
- }
+ lifecycles.add(endpoint);
+ if (endpoint.isAutoStartup()) {
+ endpoint.start();
}
if (result instanceof ApplicationListener) {
listeners.add((ApplicationListener) result);
@@ -161,6 +154,9 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean
private boolean shouldCreateEndpoint(Annotation annotation) {
Object inputChannel = AnnotationUtils.getValue(annotation, "inputChannel");
+ if (inputChannel == null && annotation instanceof InboundChannelAdapter) {
+ inputChannel = AnnotationUtils.getValue(annotation);
+ }
return (inputChannel != null && inputChannel instanceof String
&& StringUtils.hasText((String) inputChannel));
}
@@ -233,4 +229,5 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean
}
this.running = false;
}
+
}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java
index 8b7cebcc6f..800a80ea97 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java
@@ -26,6 +26,8 @@ import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
+import java.util.Date;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import org.hamcrest.Matchers;
@@ -48,6 +50,7 @@ import org.springframework.core.convert.converter.Converter;
import org.springframework.core.serializer.support.SerializingConverter;
import org.springframework.integration.annotation.Gateway;
import org.springframework.integration.annotation.GatewayHeader;
+import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.IntegrationComponentScan;
import org.springframework.integration.annotation.MessageEndpoint;
import org.springframework.integration.annotation.MessagingGateway;
@@ -66,6 +69,7 @@ import org.springframework.integration.config.EnableMessageHistory;
import org.springframework.integration.config.EnablePublisher;
import org.springframework.integration.config.GlobalChannelInterceptor;
import org.springframework.integration.config.IntegrationConverter;
+import org.springframework.integration.endpoint.MethodInvokingMessageSource;
import org.springframework.integration.endpoint.PollingConsumer;
import org.springframework.integration.history.MessageHistory;
import org.springframework.integration.history.MessageHistoryConfigurer;
@@ -80,6 +84,7 @@ import org.springframework.messaging.support.ChannelInterceptor;
import org.springframework.messaging.support.ChannelInterceptorAdapter;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.scheduling.Trigger;
+import org.springframework.scheduling.TriggerContext;
import org.springframework.scheduling.support.CronTrigger;
import org.springframework.scheduling.support.PeriodicTrigger;
import org.springframework.stereotype.Component;
@@ -94,6 +99,7 @@ import org.springframework.test.context.support.AnnotationConfigContextLoader;
*/
@ContextConfiguration(loader = AnnotationConfigContextLoader.class, classes = {EnableIntegrationTests.ContextConfiguration.class, EnableIntegrationTests.ContextConfiguration2.class})
@RunWith(SpringJUnit4ClassRunner.class)
+@DirtiesContext
public class EnableIntegrationTests {
@Autowired
@@ -159,6 +165,15 @@ public class EnableIntegrationTests {
@Autowired
private MessageChannel bytesChannel;
+ @Autowired
+ private PollableChannel counterChannel;
+
+ @Autowired
+ private PollableChannel fooChannel;
+
+ @Autowired
+ private PollableChannel messageChannel;
+
@Test
public void testAnnotatedServiceActivator() {
assertEquals(10L, TestUtils.getPropertyValue(this.serviceActivatorEndpoint, "maxMessagesPerPoll"));
@@ -224,6 +239,27 @@ public class EnableIntegrationTests {
assertNull(this.wireTapChannel.receive(0));
assertThat(this.testChannelInterceptor.getInvoked(), Matchers.greaterThan(0));
assertThat(this.fbInterceptorCounter.get(), Matchers.greaterThan(0));
+
+ assertTrue(this.context.containsBean("enableIntegrationTests.AnnotationTestService.count.inboundChannelAdapter.source"));
+ Object messageSource = this.context.getBean("enableIntegrationTests.AnnotationTestService.count.inboundChannelAdapter.source");
+ assertThat(messageSource, Matchers.instanceOf(MethodInvokingMessageSource.class));
+
+ for (int i = 0; i < 10; i++) {
+ Message> message = this.counterChannel.receive(1000);
+ assertNotNull(message);
+ assertEquals(i + 1, message.getPayload());
+ }
+
+ Message> message = this.fooChannel.receive(1000);
+ assertNotNull(message);
+ assertEquals("foo", message.getPayload());
+ assertNull(this.fooChannel.receive(10));
+
+ message = this.messageChannel.receive(1000);
+ assertNotNull(message);
+ assertEquals("bar", message.getPayload());
+ assertTrue(message.getHeaders().containsKey("foo"));
+ assertEquals("FOO", message.getHeaders().get("foo"));
}
@Test
@@ -322,6 +358,19 @@ public class EnableIntegrationTests {
return new PeriodicTrigger(1000L);
}
+ @Bean
+ public Trigger onlyOnceTrigger() {
+ return new Trigger() {
+
+ private final AtomicBoolean invoked = new AtomicBoolean();
+
+ @Override
+ public Date nextExecutionTime(TriggerContext triggerContext) {
+ return this.invoked.getAndSet(true) ? null : new Date();
+ }
+ };
+ }
+
@Bean
public PollableChannel output() {
return new QueueChannel();
@@ -427,6 +476,21 @@ public class EnableIntegrationTests {
return new QueueChannel();
}
+ @Bean
+ public PollableChannel counterChannel() {
+ return new QueueChannel();
+ }
+
+ @Bean
+ public PollableChannel fooChannel() {
+ return new QueueChannel();
+ }
+
+ @Bean
+ public PollableChannel messageChannel() {
+ return new QueueChannel();
+ }
+
@Bean
public QueueChannel numberChannel() {
QueueChannel channel = new QueueChannel();
@@ -484,6 +548,8 @@ public class EnableIntegrationTests {
@MessageEndpoint
public static class AnnotationTestService {
+ private AtomicInteger counter = new AtomicInteger();
+
@ServiceActivator(inputChannel = "input", outputChannel = "output",
poller = @Poller(maxMessagesPerPoll = "${poller.maxMessagesPerPoll}", fixedDelay = "${poller.interval}"))
@Publisher
@@ -541,6 +607,37 @@ public class EnableIntegrationTests {
assertEquals("echo", message.getHeaders().get("calledMethod"));
return this.handle(message.getPayload());
}
+
+ @InboundChannelAdapter("counterChannel")
+ public Integer count() {
+ return this.counter.incrementAndGet();
+ }
+
+ @InboundChannelAdapter(value = "fooChannel", poller = @Poller(trigger = "onlyOnceTrigger", maxMessagesPerPoll = "1"))
+ public String foo() {
+ return "foo";
+ }
+
+ @InboundChannelAdapter(value = "messageChannel", poller = @Poller(fixedDelay = "${poller.interval}", maxMessagesPerPoll = "1"))
+ public Message> message() {
+ return MessageBuilder.withPayload("bar").setHeader("foo", "FOO").build();
+ }
+
+ /*
+ * This is an error because 'InboundChannelAdapter' method must not have any arguments.
+ */
+ /*@InboundChannelAdapter("errorChannel")
+ public String error1(Object arg) {
+ return "foo";
+ }*/
+
+ /*
+ * This is an error because 'InboundChannelAdapter' return type must not be 'void'.
+ */
+ /*@InboundChannelAdapter("errorChannel")
+ public void error2() {
+ }*/
+
}
@MessagingGateway(defaultRequestChannel = "gatewayChannel", defaultHeaders = @GatewayHeader(name = "foo", value = "FOO"))
diff --git a/src/reference/docbook/configuration.xml b/src/reference/docbook/configuration.xml
index b6bb75543f..0036aaa8cb 100644
--- a/src/reference/docbook/configuration.xml
+++ b/src/reference/docbook/configuration.xml
@@ -230,6 +230,7 @@ public class FooService {
@ServiceActivator
@Splitter
@Transformer
+ @InboundChannelAdapter
The behavior of each is described in its own chapter or section within
@@ -385,6 +386,37 @@ public PollerMetadata myPoller() {
}
}]]>
+
+ @InboundChannelAdapter
+
+
+ Starting with version 4.0, the
+ @InboundChannelAdapter method annotation is available.
+ This produces a SourcePollingChannelAdapter
+ integration component based on a MethodInvokingMessageSource for the annotated method.
+ This annotation is an analogue of <int:inbound-channel-adapter>
+ XML component and has the same restrictions:
+ the method cannot have parameters, and the return type must not be void.
+ It has two attributes: value -
+ the required MessageChannel bean name and poller - an optional
+ @Poller annotation, as described above. If there
+ is need to provide some MessageHeaders, use a
+ Message<?> return type and
+ build the Message<?> within the method
+ using a MessageBuilder to
+ configure its MessageHeaders.
+
+ The first example requires that the default poller has been declared elsewhere in the application
+ context.
+
Also see .
diff --git a/src/reference/docbook/whats-new.xml b/src/reference/docbook/whats-new.xml
index c043b1eb93..9a5da506fb 100644
--- a/src/reference/docbook/whats-new.xml
+++ b/src/reference/docbook/whats-new.xml
@@ -154,6 +154,14 @@
For more information, see .
+
+ @InboundChannelAdapter
+
+ The @InboundChannelAdapter method annotation is now available.
+ It is an analogue of the <int:inbound-channel-adapter> XML component.
+ For more information, see .
+
+