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

Spring WebFlux操作Kafka报错:No subscriptions have been created求助

解决Reactor Kafka "No subscriptions have been created" 错误

问题原因

创建ReactiveKafkaConsumerTemplate时,传入的ReceiverOptions未指定要订阅的Kafka Topic,导致消费流程启动时找不到有效订阅配置,触发IllegalStateException。

修复步骤

  1. 在构建ReceiverOptions时,调用subscription()方法明确指定需要消费的Topic列表
  2. 确保消费者配置中包含必填的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 20:47:20