如何修改Kafka Flux实现低速消息接收?现有代码报Scheduler异常
尝试修改监听Kafka Topic的Kafka Flux,使其以较低速率接收消息,但当前代码抛出异常。
代码实现:
public class SampleConsumer { private static final Logger log = LoggerFactory.getLogger(SampleConsumer.class.getName()); private static final String BOOTSTRAP_SERVERS = "localhost:9092"; boolean slow = false; private static final String TOPIC = "demo-topic"; private final ReceiverOptions<Integer, String> receiverOptions; private final DateTimeFormatter dateFormat; public SampleConsumer(String bootstrapServers) { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.CLIENT_ID_CONFIG, "sample-consumer"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "sample-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); receiverOptions = ReceiverOptions.create(props); dateFormat = DateTimeFormatter.ofPattern("HH:mm:ss:SSS z dd MMM yyyy"); } public void consumeMessages(String topic, CountDownLatch latch) { ReceiverOptions<Integer, String> options = receiverOptions.subscription(Collections.singleton(topic)) .addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions)) .addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions)); Flux<ReceiverRecord<Integer, String>> kafkaFlux = KafkaReceiver.create(options).receive(); Disposable disposable = kafkaFlux.doOnNext(record -> System.out.println("Received message: value="+ record.value()) ) .subscribe(record -> { ReceiverOffset offset = record.receiverOffset(); offset.acknowledge(); }); if(slow) // logic has been implemented for this { disposable.dispose(); disposable = kafkaFlux.take(1, true).doOnNext(record -> System.out.println("Received slow message: value="+ record.value()) ) .subscribe(record -> { ReceiverOffset offset = record.receiverOffset(); offset.acknowledge(); }); } }}
异常信息:
Caused by reactor.core.Exceptions$StaticRejectedExecutionException : Scheduler unavailable
解决方案
问题根源
同一个KafkaReceiver生成的Flux实例被多次订阅,且第一次订阅后调用dispose()关闭了底层调度器,第二次订阅时调度器已不可用,导致异常。
实现方式1:根据slow标志分支处理
直接根据slow状态创建对应的处理流程,避免复用已被dispose的Flux:
import java.time.Duration; import reactor.core.publisher.Mono; public void consumeMessages(String topic, CountDownLatch latch) { ReceiverOptions<Integer, String> options = receiverOptions.subscription(Collections.singleton(topic)) .addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions)) .addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions)); Flux<ReceiverRecord<Integer, String>> kafkaFlux = KafkaReceiver.create(options).receive(); Disposable disposable; if(slow) { // 低速模式:每条消息处理后延迟1秒,控制接收速率 disposable = kafkaFlux.concatMap(record -> { System.out.println("Received slow message: value="+ record.value()); record.receiverOffset().acknowledge(); return Mono.delay(Duration.ofSeconds(1)); }) .subscribe(); } else { // 正常模式 disposable = kafkaFlux.doOnNext(record -> System.out.println("Received message: value="+ record.value())) .subscribe(record -> record.receiverOffset().acknowledge()); } }
实现方式2:多播Flux实现灵活速率控制
如果需要后续动态调整速率,可将Flux转为多播模式,避免重复创建Kafka Receiver:
import reactor.core.publisher.ConnectableFlux; import java.time.Duration; public void consumeMessages(String topic, CountDownLatch latch) { ReceiverOptions<Integer, String> options = receiverOptions.subscription(Collections.singleton(topic)) .addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions)) .addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions)); // 将Flux转为可多播的ConnectableFlux ConnectableFlux<ReceiverRecord<Integer, String>> connectableFlux = KafkaReceiver.create(options).receive().publish(); Disposable disposable; if(slow) { // 低速模式:限制每秒处理1条 disposable = connectableFlux.delayElements(Duration.ofSeconds(1)) .doOnNext(record -> System.out.println("Received slow message: value="+ record.value())) .subscribe(record -> record.receiverOffset().acknowledge()); } else { // 正常模式 disposable = connectableFlux.doOnNext(record -> System.out.println("Received message: value="+ record.value())) .subscribe(record -> record.receiverOffset().acknowledge()); } // 启动多播流 connectableFlux.connect(); }
关键注意事项
- 不要重复订阅同一个
KafkaReceiver.create(options).receive()返回的Flux,每个实例绑定唯一的Receiver生命周期,dispose后无法复用。 - 低速接收的核心是控制消息处理速率,而非频繁启停订阅,否则会破坏Kafka Receiver的稳定运行。
- 若需动态切换速率,可结合
switchOnNext操作符实现流的动态切换,确保Receiver生命周期不受影响。
内容的提问来源于stack exchange,提问作者shiva kumar
相关产品推荐
相关产品推荐

