Spring Integration XML转Java DSL:入站通道适配器方法调用
Spring Integration XML转Java DSL解决方案
原XML配置
<int:inbound-channel-adapter id="tasksForResendingToAggregatorPoller" ref="taskProcessor" method="getTaskIdForResendingToAggregator" channel="resultAggregationChannel" auto-startup="false"> <int:poller max-messages-per-poll="-1" fixed-delay="#{1 * T(org.apache.commons.lang3.time.DateUtils).MILLIS_PER_HOUR}" receive-timeout="-1"/> </int:inbound-channel-adapter> <int:channel id="resultAggregationChannel" datatype="java.lang.Long"> <int:queue capacity="${maxNumberOfMessagesInBuffers}"/> </int:channel>
现有未完成的Java DSL代码
@Bean(name= PollerMetadata.DEFAULT_POLLER) public PollerMetadata defaultPoller() { return Pollers.fixedDelay(DateUtils.MILLIS_PER_HOUR).receiveTimeout(-1).get(); } @Bean public MessageChannel resultAggregationChannel() { return MessageChannels.queue(bceMaxNumberOfMessagesInBuffers).get(); } @Bean public IntegrationFlow taskAgregator() { return IntegrationFlows.from("resultAggregationChannel") .handle(getEnrichmentTaskIdForResendingToAggregator) .get(); };
完整正确的Java DSL实现
import org.apache.commons.lang3.time.DateUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.dsl.MessageChannels; import org.springframework.integration.dsl.Pollers; import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.messaging.MessageChannel; // 注入你的TaskProcessor实例 @Autowired private TaskProcessor taskProcessor; @Bean(name = PollerMetadata.DEFAULT_POLLER) public PollerMetadata defaultPoller() { return Pollers.fixedDelay(DateUtils.MILLIS_PER_HOUR) .maxMessagesPerPoll(-1) .receiveTimeout(-1) .get(); } @Bean public MessageChannel resultAggregationChannel() { return MessageChannels.queue(bceMaxNumberOfMessagesInBuffers) .datatype(Long.class) .get(); } @Bean public IntegrationFlow tasksForResendingToAggregatorPollerFlow() { return IntegrationFlows.from(() -> taskProcessor.getTaskIdForResendingToAggregator(), spec -> spec.id("tasksForResendingToAggregatorPoller") .autoStartup(false) .poller(Pollers.fixedDelay(DateUtils.MILLIS_PER_HOUR) .maxMessagesPerPoll(-1) .receiveTimeout(-1))) .channel("resultAggregationChannel") .get(); }
关键说明
- ref和method的对应实现:通过
IntegrationFlows.from的方法引用() -> taskProcessor.getTaskIdForResendingToAggregator,直接对应XML中ref="taskProcessor"和method="getTaskIdForResendingToAggregator"的功能,轮询调用指定Bean的方法并将返回值作为消息Payload发送。 - 通道配置补全:原XML中channel指定了
datatype="java.lang.Long",在Java DSL中通过.datatype(Long.class)补全该限制。 - Poller配置补全:原XML的poller设置了
max-messages-per-poll="-1",在Pollers构建器中添加.maxMessagesPerPoll(-1)来对应。 - 适配器属性映射:通过
spec.id("tasksForResendingToAggregatorPoller")和.autoStartup(false),分别对应XML中inbound-channel-adapter的id和auto-startup属性。
内容的提问来源于stack exchange,提问作者Shayuh
相关产品推荐
相关产品推荐

