如何用IntegrationFlow实现Spring Cloud Stream Kafka消息发送并解决订阅者缺失异常
问题描述
原基于Spring Cloud Stream Kafka的Supplier代码可正常运行:
@SpringBootApplication public class DemoApplication { public static void main(String[] args) { SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic"); } @Bean Supplier<String> outbound() { return () -> { return LocalTime.now().toString(); }; } }
尝试改用IntegrationFlow实现相同功能并添加转换器时,抛出异常MessageDispatchingException: Dispatcher has no subscribers,用户代码如下:
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.Pollers; import java.time.LocalTime; @SpringBootApplication public class DemoApplication { public static void main(String[] args) { SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic"); } @Bean IntegrationFlow myFlow() { return IntegrationFlow.fromSupplier(this::myPoller, p -> p.poller(Pollers.fixedDelay(5000))) .transform(m -> { System.out.println("my transformer"); return m; }) .channel("outbound") .get(); } String myPoller() { return LocalTime.now() + " value"; } }
Gradle配置:
plugins { id 'java' id 'org.springframework.boot' version '3.2.9' id 'io.spring.dependency-management' version '1.1.6' } group = 'com.example' version = '0.0.1-SNAPSHOT' java { toolchain { languageVersion = JavaLanguageVersion.of(21) } } repositories { mavenCentral() } ext { set('springCloudVersion', "2023.0.3") } dependencies { implementation 'org.springframework.cloud:spring-cloud-stream' implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka' implementation 'org.springframework.kafka:spring-kafka' testImplementation 'org.springframework.boot:spring-boot-starter-test' testImplementation 'org.springframework.cloud:spring-cloud-stream-test-binder' testImplementation 'org.springframework.kafka:spring-kafka-test' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' } dependencyManagement { imports { mavenBom "org.springframework.cloud:spring-cloud-dependencies:${springCloudVersion}" } } tasks.named('test') { useJUnitPlatform() }
错误原因
代码中.channel("outbound")指定的是普通Spring Integration通道,该通道未与Kafka出站绑定关联,因此没有订阅者,导致消息无法分发。而原Supplier方式是通过Spring Cloud Stream的函数式绑定自动创建了与Kafka关联的出站通道(outbound-out-0),消息能正常发送到Kafka。
解决方案
以下两种方式均可实现需求,将IntegrationFlow的输出连接到Spring Cloud Stream的Kafka出站绑定:
方式1:直接绑定到自动创建的出站通道
利用配置中指定的outbound-out-0通道名(对应spring.cloud.stream.bindings.outbound-out-0.destination=out-topic),将IntegrationFlow的输出指向该通道:
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.Pollers; import java.time.LocalTime; @SpringBootApplication public class DemoApplication { public static void main(String[] args) { SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic"); } @Bean IntegrationFlow myFlow() { return IntegrationFlow.fromSupplier(this::myPoller, p -> p.poller(Pollers.fixedDelay(5000))) .transform(m -> { System.out.println("my transformer"); return m; }) .channel("outbound-out-0") // 替换为Spring Cloud Stream创建的出站通道 .get(); } String myPoller() { return LocalTime.now() + " value"; } }
方式2:使用StreamBridge发送(更灵活)
通过StreamBridge直接发送消息到指定绑定目标,无需依赖固定通道名:
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.Pollers; import java.time.LocalTime; @SpringBootApplication public class DemoApplication { public static void main(String[] args) { SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic"); } @Bean IntegrationFlow myFlow(StreamBridge streamBridge) { return IntegrationFlow.fromSupplier(this::myPoller, p -> p.poller(Pollers.fixedDelay(5000))) .transform(m -> { System.out.println("my transformer"); return m; }) .handle(message -> streamBridge.send("outbound-out-0", message.getPayload())) // 通过StreamBridge发送消息 .get(); } String myPoller() { return LocalTime.now() + " value"; } }
说明
- 方式1中,
outbound-out-0是Spring Cloud Stream根据绑定配置自动创建的出站通道,已与Kafka绑定器关联,消息会被发送到配置的out-topic。 - 方式2的
StreamBridge支持动态指定发送目标,适合需要根据业务逻辑切换主题的场景,同样会利用Spring Cloud Stream的绑定配置完成Kafka消息发送。
内容的提问来源于stack exchange,提问作者mike27
相关产品推荐
相关产品推荐

