Spring Integration如何让Service Activator多线程轮询Pollable Channel并行处理
解决方案
一、Service Activator多线程并行配置方法
你可以通过给轮询器绑定自定义线程池的方式实现指定线程数的并行处理,操作步骤如下:
- 先定义指定并行数的线程池Bean
import org.springframework.context.annotation.Bean; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.TaskExecutor; @Bean public TaskExecutor flightProcessTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 此处设置核心线程数,即你需要的并行处理最大线程数 executor.setCorePoolSize(8); executor.setMaxPoolSize(8); // 队列容量设为0,线程满了直接拒绝,避免任务堆积 executor.setQueueCapacity(0); executor.setThreadNamePrefix("flight-process-"); executor.initialize(); return executor; }
- 修改Service Activator的@Poller配置,绑定上述线程池
@ServiceActivator(inputChannel = MessageChannelConstants.get_tt_flights, poller = @Poller(fixedRate = "5000", taskExecutor = "flightProcessTaskExecutor", maxMessagesPerPoll = "5")) public void processFlightData(Message<?> receive) throws PacPlusException { // 保留原有处理逻辑即可 long startTime = System.currentTimeMillis(); }
说明:maxMessagesPerPoll配置为5是匹配你SQS适配器每次最多拉取5条消息的配置,你可以根据实际情况调整数值。
二、额外疑问解答
1. 两个指向同一Pollable Channel的Service Activator能否正常运行?
可以正常运行。QueueChannel本身是线程安全的实现,多个消费者轮询同一个队列通道时,每条消息只会被其中一个消费者拿到,不会出现重复消费的问题。
2. 这种多Service Activator的并行方式是否合理?
不合理,不推荐使用,原因如下:
- 扩展性极差:如果要调整并行数,需要新增/删除Service Activator代码,远不如直接修改线程池参数方便
- 额外资源消耗:多个Service Activator会各自独立触发轮询,容易出现大量空轮询,浪费CPU资源
- 维护成本高:多个轮询器的参数需要单独配置,很容易出现配置不一致的问题
内容的提问来源于stack exchange,提问作者user2393260
相关产品推荐
相关产品推荐

