Spring Integration DSL 关闭策略:流关闭时正在处理的文件是否会正常完成?
问题描述
我当前有如下流:
package com.example.demo.flow; import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.*; import org.springframework.integration.dsl.channel.MessageChannels; import org.springframework.integration.file.dsl.Files; import org.springframework.stereotype.Component; import java.io.File; import java.util.concurrent.Executors; /** * Created by on 03/01/2020. */ @Component @Slf4j public class TestFlow { @Bean public StandardIntegrationFlow errorChannelHandler() { return IntegrationFlows.from("testChannel") .handle(o -> { log.info("Handling error....{}", o); }).get(); } @Bean public IntegrationFlow testFile() { IntegrationFlowBuilder testChannel = IntegrationFlows.from(Files.inboundAdapter(new File("d:/input-files/")), e -> e.poller(Pollers.fixedDelay(5000L).maxMessagesPerPoll(5) .errorChannel("testChannel"))) .channel(MessageChannels.executor(Executors.newFixedThreadPool(5))) .transform(o -> { throw new RuntimeException("Failing on purpose"); }).handle(o -> { }); return testChannel.get(); } }
我计划在服务关停时停止该集成流,但此时可能有部分文件正处于流的处理流程中。请问如果我关闭该集成流,正在处理的文件会继续执行直至处理完成,还是会被直接终止对应线程、中断处理流程?
解答
核心结论
默认情况下你关闭集成流时,正在处理的文件会继续执行到完成,不会被主动中断,仅停止拉取新的文件进入处理流程。
具体行为说明
- 集成流停止的第一步是关闭入站适配器,你这里用的文件轮询器会立刻停止拉取新的文件,不会有新的任务进入处理链路。
- 已经被轮询拉取、提交到线程池的处理任务,Spring Integration不会主动调用
interrupt()方法中断线程,只要JVM还未退出,任务会一直运行到处理完成或者抛出异常终止。 - 你当前用的
Executors.newFixedThreadPool(5)生成的是普通非守护线程,只要线程未运行结束,会阻止JVM退出,直到所有任务跑完。
现有代码的潜在问题
你当前的实现有两个风险点可能导致任务被强制终止:
- 自定义线程池未交给Spring管理,Spring容器优雅停机时不会主动等待该线程池的任务执行完成,若容器销毁Bean的速度快于任务处理速度,可能出现任务运行过程中依赖的Bean被销毁,导致处理异常中断。
- 没有设置停机等待超时,若任务处理耗时超过Spring容器默认的30s优雅停机等待时长,JVM会强制退出,未完成的任务会被直接终止。
优化方案
建议将线程池交给Spring托管,开启停机任务等待配置:
@Bean public ThreadPoolTaskExecutor flowExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(5); // 停机时等待所有已提交任务执行完成 executor.setWaitForTasksToCompleteOnShutdown(true); // 最多等待60s,超时后强制退出 executor.setAwaitTerminationSeconds(60); executor.setThreadNamePrefix("flow-handler-"); executor.initialize(); return executor; }
然后在集成流中引用该线程池Bean:
.channel(MessageChannels.executor(flowExecutor()))
内容的提问来源于stack exchange,提问作者Makky
相关产品推荐
相关产品推荐

