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

如何使用Reactor Kafka消费到分区最新偏移量后停止消费者

解决方案:从Kafka主题起始位置消费至最新偏移量后自动停止

要实现从主题起始位置消费到每个分区的最新偏移量后停止监听,你需要跟踪每个分区的最新偏移量(End Offset),并在消费到该分区的最后一条消息时触发停止逻辑。你考虑的position()方法并不适用——它返回的是消费者下一个要读取的偏移量,而非分区的最新偏移量,正确的做法是获取每个分区的endOffsets。

具体实现步骤

  1. 在分区分配时获取每个分区的最新偏移量:利用addAssignListener,在分区分配完成后,获取每个分区的最新偏移量并保存。
  2. 跟踪每个分区的已消费最大偏移量:用线程安全的集合记录每个分区消费到的最大偏移量。
  3. 判断是否所有分区都已消费完毕:每次消费完一条消息后,检查当前分区的已消费偏移量是否等于该分区的最新偏移量减1(因为偏移量从0开始,最新的可消费消息偏移量是endOffset - 1),当所有分区都满足这个条件时,停止消费者。

修改后的代码示例

// 用于跟踪每个分区的最新偏移量
private final ConcurrentMap<TopicPartition, Long> partitionEndOffsets = new ConcurrentHashMap<>();
// 用于跟踪每个分区已消费的最大偏移量
private final ConcurrentMap<TopicPartition, Long> partitionConsumedOffsets = new ConcurrentHashMap<>();

ReceiverOptions<String, byte[]> receiverOptions = this.receiverOptions
        .subscription(Collections.singleton(topicName))
        .addAssignListener(partitions -> {
            log.debug("starting from beginning offset");
            // 先将消费者定位到分区起始位置
            partitions.forEach(ReceiverPartition::seekToBeginning);
            
            // 获取每个分区的最新偏移量
            KafkaConsumer<String, byte[]> consumer = partitions.get(0).consumer();
            Map<TopicPartition, Long> endOffsets = consumer.endOffsets(partitions.stream()
                    .map(p -> new TopicPartition(p.topic(), p.partition()))
                    .collect(Collectors.toList()));
            partitionEndOffsets.putAll(endOffsets);
            log.debug("got end offsets for partitions: {}", endOffsets);
        })
        .addRevokeListener(partitions -> {
            log.debug("partition revoked = {}", partitions);
            // 分区回收时清空跟踪数据
            partitions.forEach(p -> {
                TopicPartition tp = new TopicPartition(p.topic(), p.partition());
                partitionEndOffsets.remove(tp);
                partitionConsumedOffsets.remove(tp);
            });
        });

Flux<ReceiverRecord<String, byte[]>> fluxKafka = KafkaReceiver.create(receiverOptions).receive();

disposable = fluxKafka.publishOn(Schedulers.fromExecutor(scheduledExecutor))
        .subscribe(record -> {
            log.debug("received record [offset = {}, message = {}, timestamp = {}]", record.offset(),
                    new String(record.value()), record.timestamp());

            T message = avroDeserializer.deserialize(record.value());
            Optional.ofNullable(onMessageHandler)
                    .ifPresent(handler -> handler.accept(message));

            record.receiverOffset().acknowledge();

            // 更新当前分区的已消费最大偏移量
            TopicPartition tp = new TopicPartition(record.topic(), record.partition());
            partitionConsumedOffsets.put(tp, record.offset());

            // 检查是否所有分区都已消费到最新偏移量
            boolean allPartitionsCompleted = partitionEndOffsets.entrySet().stream()
                    .allMatch(entry -> {
                        Long consumedOffset = partitionConsumedOffsets.get(entry.getKey());
                        // 最新可消费偏移量是endOffset - 1
                        return consumedOffset != null && consumedOffset.equals(entry.getValue() - 1);
                    });

            if (allPartitionsCompleted) {
                log.debug("all partitions have been consumed to latest offset, stopping consumer");
                disposable.dispose();
            }
        }, throwable -> {
            log.error("error occurred while consuming", throwable);
            disposable.dispose();
        });

关键细节说明

  • 获取End Offset:通过consumer.endOffsets()方法获取分区的最新偏移量,这个值是分区中下一条待写入消息的偏移量,因此最后一条可消费消息的偏移量是endOffset - 1。
  • 线程安全集合:使用ConcurrentHashMap保证多线程环境下的安全访问,因为Reactor Kafka的消费逻辑可能在多个线程中执行。
  • 停止时机:当所有分区的已消费偏移量都达到endOffset - 1时,调用disposable.dispose()停止消费者,同时处理消费异常的情况,避免消费者挂起。

内容的提问来源于stack exchange,提问作者Rico Sancho Abarro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 21:45:35