You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.24 16:27:24