Spring WebFlux操作Kafka报错:No subscriptions have been created求助
解决Reactor Kafka "No subscriptions have been created" 错误
问题原因
创建ReactiveKafkaConsumerTemplate时,传入的ReceiverOptions未指定要订阅的Kafka Topic,导致消费流程启动时找不到有效订阅配置,触发IllegalStateException。
修复步骤
- 在构建
ReceiverOptions时,调用subscription()方法明确指定需要消费的Topic列表 - 确保消费者配置中包含必填的
group-id,并根据需求设置auto-offset-reset(如需消费历史数据可设为earliest)
修改后的代码示例
1. 更新ReactiveKafkaConsumerTemplate Bean配置
import java.util.Collections; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.KafkaProperties; import reactor.kafka.receiver.ReceiverOptions; import reactor.kafka.receiver.ReactiveKafkaConsumerTemplate; @Bean public ReactiveKafkaConsumerTemplate<String, String> reactiveKafkaConsumerTemplate( KafkaProperties properties) { Map<String, Object> props = properties.buildConsumerProperties(); // 替换为你实际的Kafka Topic名称 ReceiverOptions<String, String> receiverOptions = ReceiverOptions.create(props) .subscription(Collections.singletonList("your-target-topic")); return new ReactiveKafkaConsumerTemplate<>(receiverOptions); }
2. 调整Controller消费方法
import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Flux; import reactor.kafka.receiver.ReactiveKafkaConsumerTemplate; @RestController public class KafkaStreamController { private final ReactiveKafkaConsumerTemplate<String, String> consumerTemplate; public KafkaStreamController(ReactiveKafkaConsumerTemplate<String, String> consumerTemplate) { this.consumerTemplate = consumerTemplate; } @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> streamKafkaMessages() { return consumerTemplate.receiveAtMostOnce() .map(record -> record.value()) .doOnNext(System.out::println); } }
3. 补充配置文件示例(application.properties)
spring.kafka.consumer.group-id=webflux-kafka-consumer-group spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.bootstrap-servers=localhost:9092
验证要点
- 确认Kafka集群可正常访问
- 目标Topic已存在且有消息产生
- 消费者groupId未被其他消费进程占用(或根据需求调整)
内容的提问来源于stack exchange,提问作者user3878073
相关产品推荐
相关产品推荐

