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

Spring Integration IntegrationFlowAdapter如何使用?求可用示例

适配最新版Spring Integration的MyFlowAdapter示例

原官方文档示例在最新版Spring Integration中无法编译,主要是API变更导致的,以下是修正后的可运行版本:

@Component
public class MyFlowAdapter extends IntegrationFlowAdapter {

    private final AtomicBoolean invoked = new AtomicBoolean();

    public Date nextExecutionTime(TriggerContext triggerContext) {
        return this.invoked.getAndSet(true) ? null : new Date();
    }

    @Override
    protected IntegrationFlowDefinition<?> buildFlow() {
        return from(this::messageSource,
                        e -> e.poller(p -> p.trigger(this::nextExecutionTime)))
                .split(this)
                .transform(this)
                .aggregate(a -> a.processor(this))
                .enrichHeaders(Collections.singletonMap("thing1", "THING1"))
                .filter(this)
                .handle(this)
                .channel(c -> c.queue("myFlowAdapterOutput"));
    }

    public String messageSource() {
        return "T,H,I,N,G,2";
    }

    @Splitter
    public String[] split(String payload) {
        return StringUtils.commaDelimitedListToStringArray(payload);
    }

    @Transformer
    public String transform(String payload) {
        return payload.toLowerCase();
    }

    @Aggregator
    public String aggregate(List<String> payloads) {
        return payloads.stream().collect(Collectors.joining());
    }

    @Filter
    public boolean filter(@Header(required = false) String thing1) {
        return thing1 != null;
    }

    @ServiceActivator
    public String handle(String payload, @Header String thing1) {
        return payload + ":" + thing1;
    }

}

关键修改说明:

  • aggregate方法调整:原代码中aggregate(a -> a.processor(this, null), null)的写法已被废弃,最新版只需传入a -> a.processor(this)即可,框架会自动识别@Aggregator注解的方法
  • Filter参数修正:原代码中@Header Optional<String> thing1的写法在最新版Spring中可能引发类型转换问题,改用@Header(required = false) String thing1更直接,判断非空即可
  • 依赖导入确认:确保项目中已正确引入Spring Integration最新稳定版依赖(如6.x系列),同时导入以下必要类:
    import java.util.Date;
    import java.util.List;
    import java.util.Collections;
    import java.util.concurrent.atomic.AtomicBoolean;
    import org.springframework.integration.dsl.IntegrationFlowAdapter;
    import org.springframework.integration.dsl.IntegrationFlowDefinition;
    import org.springframework.messaging.handler.annotation.Header;
    import org.springframework.scheduling.TriggerContext;
    import org.springframework.stereotype.Component;
    import org.springframework.util.StringUtils;
    import java.util.stream.Collectors;
    

内容的提问来源于stack exchange,提问作者Chris

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 01:31:20