如何配置PubSubInboundChannelAdapter优先级优先处理高优通道消息
实现方案
完全可以实现你要的绝对优先级消费逻辑,核心是从线程资源隔离和消息拉取控制两个层面入手,彻底避免低优先级通道阻塞高优先级通道的处理流程,具体实现步骤如下:
1. 第一步:给两个入站适配器分配独立的专属线程池
默认配置下PubSubInboundChannelAdapter会共用Spring全局任务执行器,这是低优先级任务占满线程、阻塞高优先级消息处理的核心原因。你需要给两个通道绑定完全隔离的线程池,从资源层面先做拆分:
// 高优先级通道专属线程池,核心/最大线程数根据你实际的消费并发需求配置 @Bean(name = "highPriorityExecutor") public TaskExecutor highPriorityExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(10); executor.setMaxPoolSize(20); executor.setThreadNamePrefix("high-priority-pubsub-"); executor.initialize(); return executor; } // 低优先级通道专属线程池,线程数配置为远小于高优先级池即可 @Bean(name = "lowPriorityExecutor") public TaskExecutor lowPriorityExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(2); executor.setMaxPoolSize(5); executor.setThreadNamePrefix("low-priority-pubsub-"); executor.initialize(); return executor; }
修改你原来的两个适配器定义,分别绑定对应的专属线程池:
@Bean public PubSubInboundChannelAdapter messageMediaChannelAdapter( @Qualifier(Channels.STORAGE_OBJECT_INPUT) MessageChannel inputChannel, @Qualifier("highPriorityExecutor") TaskExecutor highPriorityExecutor, PubSubTemplate pubSubTemplate ) { PubSubInboundChannelAdapter adapter = this.createMessageChannelAdapter(inputChannel, Channels.STORAGE_OBJECT_INPUT, pubSubTemplate); adapter.setTaskExecutor(highPriorityExecutor); // 高优先级通道可适当调大批量拉取大小,提升消费吞吐 adapter.setBatchSize(10); return adapter; } @Bean public PubSubInboundChannelAdapter messageMediaResizeChannelAdapter( @Qualifier(Channels.MEDIA_RESIZE_INPUT) MessageChannel inputChannel, @Qualifier("lowPriorityExecutor") TaskExecutor lowPriorityExecutor, PubSubTemplate pubSubTemplate ) { PubSubInboundChannelAdapter adapter = this.createMessageChannelAdapter(inputChannel, Channels.MEDIA_RESIZE_INPUT, pubSubTemplate); adapter.setTaskExecutor(lowPriorityExecutor); // 低优先级通道默认先不启动,等高优先级通道无积压时再启动拉取 adapter.setAutoStartup(false); return adapter; }
2. 第二步:根据高优先级通道积压状态控制低优先级通道的拉取动作
仅做线程隔离还达不到「高优先级消息全部处理完才处理低优先级」的要求,你需要增加一个轻量的定时检测任务,执行逻辑如下:
- 每隔1~3秒调用PubSub客户端,查询
Channels.STORAGE_OBJECT_INPUT对应订阅的未投递消息积压量 - 只要积压量大于0,或者高优先级线程池还有活跃的工作线程,就保持低优先级适配器为
stop()状态,不主动拉取任何低优先级消息 - 当检测到高优先级订阅积压为0,且高优先级线程池无活跃工作线程,持续3秒没有新的高优先级消息进入时,再调用低优先级适配器的
start()方法,恢复低优先级消息的拉取和消费 - 一旦检测到高优先级通道出现新的积压,立刻调用低优先级适配器的
stop()方法暂停消费,等所有正在处理的低优先级消息执行完(线程池任务不会被强制中断,不会出现消息丢失),所有线程资源全部留给高优先级通道使用
注意:不要用本地内存计数判断高优先级消息积压,服务重启、消息重投等场景会导致计数不准,直接调用PubSub官方客户端查询订阅的
numUndeliveredMessages监控指标是最可靠的方案。
方案效果
这套逻辑可以完全满足你的需求:
- 高优先级通道的线程资源永远不会被低优先级任务占用,新的高优先级消息进来可以立刻得到处理
- 只要高优先级通道还有未处理完的消息,低优先级通道根本不会拉取新消息到服务本地,完全不存在阻塞高优先级流程的可能
- 暂停低优先级适配器时不会强制中断正在处理的低优先级消息,不会出现消息异常丢失、重复消费概率大幅升高的问题
如果你的业务消息量不大,也可以用更轻量的实现:不需要定时检测,每次高优先级消息处理完成后就检查一次积压情况,确认无积压时手动触发低优先级通道拉取固定条数的消息,处理完立刻暂停即可。
内容的提问来源于stack exchange,提问作者Richard Lindhout
相关产品推荐
相关产品推荐

