Spring Boot集成Kafka发送消息出现Dispatcher has no subscribers错误怎么解决
问题背景
我需要实现Spring Boot应用向Kafka主题发送消息的能力,整体实现方案如下:
- 采用Spring Boot + Spring Cloud Stream Kafka集成实现
- Kafka服务通过Docker启动运行
- Spring Boot应用简单配置后连接Kafka发送消息
pom.xml依赖配置
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.5.1</version> <relativePath /> <!-- lookup parent from repository --> </parent> <dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-dependencies</artifactId> <version>2020.0.3</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-stream-binder-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-stream-binder-kafka-streams</artifactId> </dependency> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-stream-kafka</artifactId> </dependency> </dependencies>
application.yml配置
spring: kafka: bootstrap-servers: - localhost:19091 cloud: stream: bindings: kafka-log-publisher: binder: kafka destination: com.tonitingaurav.kafka.log default-binder: kafka kafka: binder: brokers: - localhost:19091
消息通道配置代码
@Configuration @EnableIntegration public class LogProducerKafkaConfig { @Bean("kafka-log-publisher") public MessageChannel kafkaLogPublisher() { return new DirectChannel(); } }
消息发布工具类
@Component public class LogEventPublisher { @Autowired @Qualifier("kafka-log-publisher") MessageChannel messageChannel; public void logMessage(Log log) { Message<Log> message = MessageBuilder.withPayload(log).build(); messageChannel.send(message); } }
测试接口代码
@RestController public class EmployeeController { private static final Logger LOGGER = LoggerFactory.getLogger(EmployeeController.class); @Autowired private LogEventPublisher logEventPublisher; @GetMapping public ResponseEntity<Employees> getAll() { LOGGER.info("Getting All Employees"); logEventPublisher.logMessage(new Log("Getting All Employees")); // 省略返回逻辑 } }
调用接口报错堆栈
org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'employee-management-1.kafka-log-publisher'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=com.tonitingaurav.microservice.audit.Log@3e170b0f, headers={id=e549b9d1-8ada-8032-2ff9-6cf1e02bae53, timestamp=1638119747020}] at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:76) ~[spring-integration-core-5.5.0.jar:5.5.0] // 省略中间堆栈 Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:139) ~[spring-integration-core-5.5.0.jar:5.5.0] at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106) ~[spring-integration-core-5.5.0.jar:5.5.0] at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72) ~[spring-integration-core-5.5.0.jar:5.5.0] ... 62 common frames omitted
错误原因与解决方案
错误原因
你手动创建的kafka-log-publisher MessageChannel没有和Spring Cloud Stream的Kafka binder绑定,Spring Cloud Stream不知道要将该通道的消息转发到Kafka集群,通道没有对应的订阅者处理消息,因此触发Dispatcher无订阅者异常。
解决方案
推荐使用Spring Cloud Stream官方提供的StreamBridge实现动态消息发送,无需自定义通道和复杂绑定配置,修改步骤如下:
- 删除手动创建
kafka-log-publisherBean的配置类LogProducerKafkaConfig - 修改
LogEventPublisher代码,注入StreamBridge发送消息:
@Component public class LogEventPublisher { @Autowired private StreamBridge streamBridge; public void logMessage(Log log) { // 第一个参数对应配置文件中binding的名称,第二个参数为消息体 streamBridge.send("kafka-log-publisher", log); } }
- 保留原有application.yml配置即可,无需修改。
如果要使用传统@EnableBinding方式(兼容旧版本),也可按如下方式修改:
- 定义输出通道接口:
public interface LogSource { @Output("kafka-log-publisher") MessageChannel output(); }
- 在启动类或配置类上添加注解
@EnableBinding(LogSource.class) - 注入
LogSource发送消息:
@Component public class LogEventPublisher { @Autowired private LogSource logSource; public void logMessage(Log log) { Message<Log> message = MessageBuilder.withPayload(log).build(); logSource.output().send(message); } }
修改完成后重启应用,调用接口即可正常将消息发送到指定Kafka主题。
内容的提问来源于stack exchange,提问作者user3082820
相关产品推荐
相关产品推荐

