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

如何修改Kafka Flux实现低速消息接收?现有代码报Scheduler异常

问题:修改Kafka Flux实现低速消息接收时抛出Scheduler unavailable异常

尝试修改监听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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 00:10:35