You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 18:32:02