使用spring-integration-mail时Dispatcher无订阅者问题排查
动态IMAP邮件接收器重启后出现"Dispatcher has no subscribers"异常解决思路
问题背景
我们有一个多租户应用,允许用户按需创建IMAP邮件抓取器(又称邮件接收器)并单独处理。选用spring-integration-mail(v6.1.4),依据官方文档采用动态注册/移除组件的方式,支持用户从UI创建/编辑/暂停/恢复邮件接收器。
核心代码
注册邮件接收器
integrationFlowContext.registration( IntegrationFlow .from( Mail.imapInboundAdapter( factory.create(imapConfiguration) // 创建ImapMailReceiver实例 ) ) { it.poller( Pollers .trigger(PeriodicTrigger(ofSeconds(10))) .taskExecutor(taskExecutor) .transactional(transactionalManager) .errorHandler(errorHandler) ) } .log<Message<*>> { mailReceiverLog.info { "Receiving message: Message(headers=${it.headers.map { header -> "${header.key}=${header.value}" }}, payload=${it.payload})" } } .channel(MessageChannels.executor(taskExecutor)) .transform(MailReceiverMessageTransformer()) .handle { payload: MailMessage -> try { // 业务处理逻辑 } catch (e: Throwable) { log.error( "Could not handle mail message: $payload, execution will continue", e ) } } .get() ) .autoStartup(true) .id(id) .useFlowIdAsPrefix() .register()
移除邮件接收器
val toRemove = integrationFlowContext.getRegistrationById(id) if (toRemove != null) { toRemove.destroy() }
配套组件配置
@Bean fun mailReceiverTaskExecutor(): TaskExecutor = ThreadPoolTaskExecutor().apply { corePoolSize = 20 } @Bean fun mailReceiverTransactionManager(): PseudoTransactionManager = PseudoTransactionManager()
异常现象
配置启动后运行稳定,但用户禁用再重新启用邮件接收器时,出现Dispatcher has no subscribers异常,无法从IMAP服务器收取邮件。日志显示通道已重新注册订阅者,但异常仍发生:
// 禁用(即移除邮件接收器)开始 Unregistering mail receiver: mailreceiver-1 stopped bean 'mailreceiver-1.org.springframework.integration.config.SourcePollingChannelAdapterFactoryBean#0' Removing {bridge} as a subscriber to the 'mailreceiver-1.channel#1' channel Channel 'backend.mailreceiver-1.channel#1' has 0 subscriber(s). stopped bean 'mailreceiver-1.org.springframework.integration.config.ConsumerEndpointFactoryBean#0' Removing {transformer} as a subscriber to the 'mailreceiver-1.channel#2' channel Channel 'backend.mailreceiver-1.channel#2' has 0 subscriber(s). stopped bean 'mailreceiver-1.org.springframework.integration.config.ConsumerEndpointFactoryBean#1' Removing {service-activator} as a subscriber to the 'mailreceiver-1.channel#3' channel Channel 'backend.mailreceiver-1.channel#3' has 0 subscriber(s). stopped bean 'mailreceiver-1.org.springframework.integration.config.ConsumerEndpointFactoryBean#2' Mail receiver successfully unregistered: mailreceiver-1 // 启用(即注册邮件接收器)开始 Registering mail receiver: mailreceiver-1 Adding {service-activator} as a subscriber to the 'mailreceiver-1.channel#3' channel Channel 'backend.mailreceiver-1.channel#3' has 1 subscriber(s). started bean 'mailreceiver-1.org.springframework.integration.config.ConsumerEndpointFactoryBean#2' Adding {transformer} as a subscriber to the 'mailreceiver-1.channel#2' channel Channel 'backend.mailreceiver-1.channel#2' has 1 subscriber(s). started bean 'mailreceiver-1.org.springframework.integration.config.ConsumerEndpointFactoryBean#1' Adding {bridge} as a subscriber to the 'mailreceiver-1.channel#1' channel Channel 'backend.mailreceiver-1.channel#1' has 1 subscriber(s). started bean 'mailreceiver-1.org.springframework.integration.config.ConsumerEndpointFactoryBean#0' started bean 'mailreceiver-1.org.springframework.integration.config.SourcePollingChannelAdapterFactoryBean#0' Mail receiver successfully registered: mailreceiver-1 org.springframework.integration.support.MessagingExceptionWrapper: null at org.springframework.integration.endpoint.AbstractPollingEndpoint.messageReceived(AbstractPollingEndpoint.java:477) at org.springframework.integration.endpoint.AbstractPollingEndpoint.doPoll(AbstractPollingEndpoint.java:460) at java.base/jdk.internal.reflect.GeneratedMethodAccessor112.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source) at java.base/java.lang.reflect.Method.invoke(Unknown Source) at org.springframework.aop.support.AopUtils.invokeJoinpointUsingReflection(AopUtils.java:343) at org.springframework.aop.framework.ReflectiveMethodInvocation.invokeJoinpoint(ReflectiveMethodInvocation.java:196) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:163) at org.springframework.transaction.interceptor.TransactionInterceptor$1.proceedWithInvocation(TransactionInterceptor.java:123) at org.springframework.transaction.interceptor.TransactionAspectSupport.invokeWithinTransaction(TransactionAspectSupport.java:391) at org.springframework.transaction.interceptor.TransactionInterceptor.invoke(TransactionInterceptor.java:119) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:184) at org.springframework.aop.framework.JdkDynamicAopProxy.invoke(JdkDynamicAopProxy.java:244) at jdk.proxy2/jdk.proxy2.$Proxy180.call(Unknown Source) at org.springframework.integration.endpoint.AbstractPollingEndpoint.pollForMessage(AbstractPollingEndpoint.java:412) at org.springframework.integration.endpoint.AbstractPollingEndpoint.lambda$createPoller$4(AbstractPollingEndpoint.java:348) at org.springframework.integration.util.ErrorHandlingTaskExecutor.lambda$execute$0(ErrorHandlingTaskExecutor.java:57) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) at java.base/java.lang.Thread.run(Unknown Source) Caused by: org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'backend.mailreceiver-1.channel#1'. at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:76) at org.springframework.integration.channel.AbstractMessageChannel.sendInternal(AbstractMessageChannel.java:375) at org.springframework.integration.channel.AbstractMessageChannel.sendWithMetrics(AbstractMessageChannel.java:346) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:326) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:299) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) at org.springframework.integration.endpoint.SourcePollingChannelAdapter.handleMessage(SourcePollingChannelAdapter.java:196) at org.springframework.integration.endpoint.AbstractPollingEndpoint.messageReceived(AbstractPollingEndpoint.java:474) ... 19 common frames omitted Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:139) at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106) at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72) ... 29 common frames omitted
排查与更新
更新1:初步怀疑
怀疑executor-channel线程不足导致该问题?
更新2:代码重构
根据建议,已重构为轻量动态流(邮件抓取器)+ 静态处理流的模式:
动态邮件抓取流
integrationFlowContext.registration( IntegrationFlow .from( Mail.imapInboundAdapter( factory.create(imapConfiguration) // 创建ImapMailReceiver实例 ) ) { it.poller( Pollers .trigger(PeriodicTrigger(ofSeconds(10))) .taskExecutor(taskExecutor) .transactional(transactionalManager) .errorHandler(errorHandler) ) } .log<Message<*>> { mailReceiverLog.info { "Receiving message: Message(headers=${it.headers.map { header -> "${header.key}=${header.value}" }}, payload=${it.payload})" } } .channel("handleMailChannel") .get() ) .autoStartup(true) .id(id) .useFlowIdAsPrefix() .register()
静态消息处理流
integrationFlowContext.registration( IntegrationFlow .from("handleMailChannel") .handle { payload: MailMessage -> try { // 业务处理逻辑 } catch (e: Throwable) { log.error( "Could not handle mail message: $payload, execution will continue", e ) } } .get() ) .autoStartup(true) .id(id) .useFlowIdAsPrefix() .register()
内容的提问来源于stack exchange,提问作者nKognito
相关产品推荐
相关产品推荐

