Spring Cloud Stream发布事件报错:Dispatcher has no subscribers for channel
问题:Spring Cloud Stream发布事件报错Dispatcher has no subscribers
我正尝试为微服务创建通用EventPublisher,代码实现如下:
EventPublisher类
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.messaging.support.MessageBuilder; @EnableBinding(TaskEventStreams.class) public class EventPublisher { private final TaskEventStreams source; @Autowired public EventPublisher(TaskEventStreams source) { this.source = source; } public void publishEvent(BaseEvent event) { publishEvent(event, event.getTenant()); } public void publishEvent(BaseEvent event, String key) { source.taskEventOutput().send( MessageBuilder .withPayload(event) .copyHeaders(event.getHeaders()) .setHeader("partitionKey", key) .build()); } }
TaskEventStreams接口
import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.Output; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.SubscribableChannel; public interface TaskEventStreams { String TASK_EVENT_OUTPUT = "task-event-output"; @Output(TASK_EVENT_OUTPUT) MessageChannel taskEventOutput(); }
应用配置(application.yml)
spring: cloud: stream: default-binder: kafka kafka: binder: brokers: localhost:9092 auto-create-topics: true bindings: task-event-output: destination: tasks group: tasks binder: kafka contentType: application/json
启动类代码
@SpringBootApplication(exclude = {DataSourceAutoConfiguration.class}) public class Main { @Autowired EventPublisher eventPublisher; public static void main(String[] args) { SpringApplication.run(Main.class, args); System.out.println("Hello world!"); } @PostConstruct public void init() { BaseEvent event = TaskCreatedEvent.builder().tenant("tenant").build(); eventPublisher.publishEvent(event); } }
异常信息
org.springframework.beans.factory.BeanCreationException: Error creating bean with name 'main': Invocation of init method failed; nested exception is org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'task-event-output'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[255], headers={id=06b61f95-24d8-7f35-3e57-9ae262c2ca47, partitionKey=tenant, contentType=application/json, tenant=tenant, timestamp=1675063077441}], failedMessage=GenericMessage [payload=byte[255], headers={id=06b61f95-24d8-7f35-3e57-9ae262c2ca47, partitionKey=tenant, contentType=application/json, tenant=tenant, timestamp=1675063077441}] at org.springframework.beans.factory.annotation.InitDestroyAnnotationBeanPostProcessor
请问是否存在无消费者则无法生产事件的限制,或是我哪里配置有误?
解决方案
核心原因
- 时机问题:
@PostConstruct在Spring容器初始化Bean阶段执行,此时Spring Cloud Stream的绑定通道还未完全初始化完成,Kafka binder还没完成和broker的连接、通道订阅,导致消息发送时通道没有订阅者。 - 配置冗余:输出通道配置
group是多余的,group属性是针对消费者输入通道的配置,生产者输出通道不需要指定group。
解决步骤
1. 修正配置文件
去掉输出通道的group配置:
spring: cloud: stream: default-binder: kafka kafka: binder: brokers: localhost:9092 auto-create-topics: true bindings: task-event-output: destination: tasks binder: kafka contentType: application/json
2. 调整事件发布时机
不要在@PostConstruct中发布事件,改用以下方式之一:
- 使用
ApplicationListener监听ApplicationReadyEvent:该事件在Spring容器完全启动、所有Bean初始化完成后触发,此时通道已准备就绪。
@SpringBootApplication(exclude = {DataSourceAutoConfiguration.class}) public class Main implements ApplicationListener<ApplicationReadyEvent> { @Autowired EventPublisher eventPublisher; public static void main(String[] args) { SpringApplication.run(Main.class, args); System.out.println("Hello world!"); } @Override public void onApplicationEvent(ApplicationReadyEvent event) { BaseEvent taskEvent = TaskCreatedEvent.builder().tenant("tenant").build(); eventPublisher.publishEvent(taskEvent); } }
- 使用
CommandLineRunner/ApplicationRunner:这两个接口的run方法会在Spring应用启动完成后执行。
@SpringBootApplication(exclude = {DataSourceAutoConfiguration.class}) public class Main implements CommandLineRunner { @Autowired EventPublisher eventPublisher; public static void main(String[] args) { SpringApplication.run(Main.class, args); System.out.println("Hello world!"); } @Override public void run(String... args) throws Exception { BaseEvent taskEvent = TaskCreatedEvent.builder().tenant("tenant").build(); eventPublisher.publishEvent(taskEvent); } }
3. 可选:配置通道为非严格模式(如果需要在启动阶段强制发送)
如果必须在早期阶段发送消息,可以配置通道的fail-on-error为false,或者设置send-timeout让消息等待通道就绪:
spring: cloud: stream: bindings: task-event-output: # 其他配置... producer: fail-on-error: false send-timeout: 5000 # 等待5秒再放弃
关键说明
- Spring Cloud Stream生产者本身不要求必须有消费者才能发送消息,报错的根本原因是发送时机过早,通道未初始化完成,而非没有消费者。
- 去掉输出通道的
group配置,避免不必要的混淆,group是消费者用来维护消费偏移量的属性,生产者不需要。
内容的提问来源于stack exchange,提问作者Tarun Lalwani
相关产品推荐
相关产品推荐

