Spring Integration中QueueChannel默认同步执行问题咨询
我在使用Spring Integration时,通过@Scheduled标注的方法,借助QueueChannel将消息发送到标注@ServiceActivator的方法。原本预期@ServiceActivator标注的方法默认会在独立线程中执行,但实际并非如此——接收方法和发送方的@Scheduled方法使用同一线程,接收方的阻塞操作会直接阻塞发送方,整个流程是同步的。
配置代码
通道配置
@Configuration public class SpecificIntegrationConfig { @Bean public MessageChannel dirListingChannel() { return new QueueChannel(5); } }
消息网关
@MessagingGateway public interface MessageGateway { @Gateway(requestChannel = "dirListingChannel") void sendListing(List<DirEntry> entries); }
消息发送逻辑(SFTP目录轮询)
@RequiredArgsConstructor public class SftpPollerComponent { private final SftpClient sftpClient; private final MessageGateway messageGateway; @Value("${remote.dir}") private final String dir; @Scheduled(fixedDelay = 2_000) public void pollRemote() throws IOException { final ArrayList<SftpClient.DirEntry> dirEntries = new ArrayList<>(20); for (final SftpClient.DirEntry entry : this.sftpClient.readDir(this.dir)) { dirEntries.add(entry); } log.info("Sender thread: {}", Thread.currentThread().getName()); this.messageGateway.sendListing(dirEntries); } }
消息接收逻辑
@Component public class DirEntryService { @ServiceActivator(inputChannel = "dirListingChannel") public void handleFiles(final List<DirEntry> files) throws InterruptedException { log.info("Receiving thread name: {}", Thread.currentThread().getName()); log.info("Sleeping.."); Thread.sleep(6000); log.info("Service awake."); } }
日志现象
从日志可以看到,发送和接收使用的是同一个scheduling-1线程,接收方的Thread.sleep(6000)直接阻塞了发送方的下一次轮询:
2024-11-28T17:45:31.503Z INFO 12987 --- [ scheduling-1] c.e.i.i.component.SftpPollerComponent : Sender thread: scheduling-1 2024-11-28T17:45:32.478Z INFO 12987 --- [ scheduling-1] c.e.i.ingest.service.DirEntryService : Receiving thread name: scheduling-1 2024-11-28T17:45:32.478Z INFO 12987 --- [ scheduling-1] c.e.i.ingest.service.DirEntryService : There are 10 files 2024-11-28T17:45:32.479Z INFO 12987 --- [ scheduling-1] c.e.i.ingest.service.DirEntryService : Sleeping.. 2024-11-28T17:45:38.479Z INFO 12987 --- [ scheduling-1] c.e.i.ingest.service.DirEntryService : Service awake. ...
疑问与官方文档参考
我原本以为指定QueueChannel后,Spring会默认自动调度轮询队列并使用其他线程。但根据Spring官方文档:
If you want the polling to be asynchronous, a poller can optionally specify a task-executor attribute that points to an existing instance of any TaskExecutor bean
这说明当前的同步执行是默认行为,但这和可轮询队列的设计理念相悖——默认情况下接收会阻塞发送方,失去了队列解耦异步的意义。
解决方案
要实现异步执行,需要给@ServiceActivator配置Poller并指定任务执行器:
方式1:使用默认TaskExecutor
直接在@ServiceActivator上添加@Poller注解,指定Spring默认提供的taskExecutor线程池:
@Component public class DirEntryService { @ServiceActivator(inputChannel = "dirListingChannel", poller = @Poller(taskExecutor = "taskExecutor")) public void handleFiles(final List<DirEntry> files) throws InterruptedException { log.info("Receiving thread name: {}", Thread.currentThread().getName()); log.info("Sleeping.."); Thread.sleep(6000); log.info("Service awake."); } }
方式2:自定义TaskExecutor
如果需要自定义线程池参数,先定义自己的TaskExecutor Bean:
@Configuration public class AsyncConfig { @Bean public TaskExecutor customTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(3); executor.setMaxPoolSize(5); executor.setQueueCapacity(10); executor.setThreadNamePrefix("dir-entry-"); executor.initialize(); return executor; } }
然后在@Poller中指定自定义Bean:
@ServiceActivator(inputChannel = "dirListingChannel", poller = @Poller(taskExecutor = "customTaskExecutor"))
方式3:配置全局默认Poller
如果多个组件需要异步处理,可配置全局默认Poller避免重复配置:
@Configuration @EnableIntegration public class IntegrationConfig { @Bean(name = PollerMetadata.DEFAULT_POLLER) public PollerMetadata defaultPoller(TaskExecutor taskExecutor) { PollerMetadata pollerMetadata = new PollerMetadata(); pollerMetadata.setTaskExecutor(taskExecutor); pollerMetadata.setTrigger(new PeriodicTrigger(100)); // 轮询间隔 return pollerMetadata; } }
绑定到QueueChannel的@ServiceActivator会自动使用该Poller实现异步。
内容的提问来源于stack exchange,提问作者D-Dᴙum

