使用SpringIntegrationTest向DirectChannel发消息报‘无订阅者’异常的解决
问题分析
你遇到的MessageDispatchingException本质是测试时,直接向mqttChannelAdapter内部创建的DirectChannel发送消息时,该通道还没有订阅者。生产环境中适配器启动后集成流会自动绑定通道,但测试里noAutoStartup阻止了适配器启动,导致集成流与通道的订阅关系未正确建立(或发送时机过早,订阅者尚未完成注册)。
可行解决方案
方案1:将输出通道抽为独立Bean(推荐)
把适配器的输出通道单独定义为Spring Bean,让集成流直接订阅这个通道。无论适配器是否启动,集成流都会在上下文初始化时完成订阅,彻底摆脱对适配器启动状态的依赖。
修改配置类:
@Configuration class MqttIntegrationConfig { // 单独定义输入通道Bean @Bean fun mqttInputChannel(): DirectChannel { return MessageChannels.direct().get() } @Bean @Qualifier("mqttChannelAdapter") fun mqttChannelAdapter(mqttInputChannel: DirectChannel): MessageProducerSupport { val adapter = MqttPahoMessageDrivenChannelAdapter(...) // ... 其他配置 adapter.outputChannel = mqttInputChannel return adapter } @Bean fun mqttInbound(mqttInputChannel: DirectChannel): IntegrationFlow { // 直接从独立的通道构建集成流 return IntegrationFlows.from(mqttInputChannel) .handle<String> { payload, headers -> logger.info("payload=$payload, headers=$headers") payload } .channel(mqttOutboundChannel()) .get() } // 其他配置省略 }
测试时两种发送方式都可行:
// 方式1:直接注入mqttInputChannel发送 @Autowired private lateinit var mqttInputChannel: DirectChannel @Test fun mytest() { mqttInputChannel.send(GenericMessage(1)) } // 方式2:通过适配器的outputChannel发送(此时指向同一个通道) @Test fun mytest() { mqttChannelAdapter.outputChannel?.send(GenericMessage(1)) }
方案2:等待通道订阅者注册完成
如果不想修改配置结构,可以在测试中添加等待逻辑,确保DirectChannel的订阅者(集成流端点)注册完成后再发送消息,利用Spring Integration的TestUtils工具类实现:
import org.springframework.integration.test.util.TestUtils @Test fun mytest() { val directChannel = mqttChannelAdapter.outputChannel as DirectChannel // 等待最多1秒,直到通道有订阅者 val hasSubscribers = TestUtils.waitForCondition( { directChannel.subscribers.isNotEmpty() }, 1000, "No subscribers registered for the DirectChannel" ) require(hasSubscribers) { "Timed out waiting for channel subscribers" } directChannel.send(GenericMessage(1)) }
方案3:直接向集成流的输入通道发送消息
通过IntegrationFlowContext获取集成流实例,直接向其输入通道发送消息,绕开适配器的状态限制:
@Autowired private lateinit var integrationFlowContext: IntegrationFlowContext @Test fun mytest() { val flowRegistration = integrationFlowContext.getRegistrationById("mqttInbound") ?: error("IntegrationFlow 'mqttInbound' not found") val inputChannel = flowRegistration.flow.inputChannel ?: error("Input channel not found for IntegrationFlow 'mqttInbound'") inputChannel.send(GenericMessage(1)) }
补充说明
- 生产环境中
DirectChannel正常工作,是因为适配器启动后集成流已完成通道订阅,且消息不会在集成流就绪前到达; QueueChannel能正常运行是因为它会缓存消息直到有订阅者消费,而DirectChannel是同步模式,要求发送时必须有订阅者存在。
内容的提问来源于stack exchange,提问作者Jay Han
相关产品推荐
相关产品推荐

