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

Reactor Kafka消费者出错停止后新建消费者的实现与测试问题

问题根因

  • 测试模拟的是receiveAutoAck()返回的Flux直接抛出全局错误,并非处理单条消费记录时产生的和元素绑定的异常,onErrorContinue仅支持捕获单元素处理阶段的关联异常,无法捕获该全局错误,会直接导致流异常终止,repeat没有触发机会。
  • retryWhen重试耗尽后抛出的RetryExhaustedException未被你定义的异常处理逻辑捕获,没有走到设置repeatConsumer为true的分支,repeat的触发条件始终不满足。
  • repeat操作符仅会在上游流正常发送complete信号时才会判断是否重新订阅,上游异常终止的场景下repeat不会生效。

代码修复方案

调整readWrite方法逻辑,区分单记录异常和全局流异常,将全局异常转换为正常结束信号触发repeat:

public Flux<String> readWrite(String destTopic) {
    // 用Flux.defer包裹,每次重新订阅都初始化新的状态变量,避免状态残留
    return Flux.defer(() -> {
        AtomicBoolean repeatConsumer = new AtomicBoolean(false);
        return kafkaConsumerTemplate
                .receiveAutoAck()
                .doOnNext(consumerRecord -> log.debug("received key={}, value={} from topic={}, offset={}",
                        consumerRecord.key(),
                        consumerRecord.value(),
                        consumerRecord.topic(),
                        consumerRecord.offset())
                )
                .doOnNext(s-> sendToKafka(s,destinationTopic))
                .map(ConsumerRecord::value)
                .doOnNext(record -> log.debug("successfully consumed {}={}", Metric[].class.getSimpleName(), record))
                // 包装所有异常为自定义类型,统一处理
                .onErrorMap(throwable -> {
                    if (throwable instanceof ReceiverRecordException) {
                        return throwable;
                    }
                    // 全局流异常无关联record,传null标记
                    return new ReceiverRecordException(null, throwable);
                })
                .doOnError(exception -> log.debug("Error occurred while processing the message, attempting retry. Error message: {}", exception.getMessage()))
                .retryWhen(Retry.backoff(Integer.parseInt(retryAttempts), Duration.ofSeconds(Integer.parseInt(retryAttemptsDelay))).transientErrors(true))
                // 处理单条记录重试耗尽的场景
                .onErrorContinue((exception, errorRecord) -> {
                    if (exception instanceof ReceiverRecordException recordException) {
                        log.debug("Retries exhausted for single record: {}", recordException);
                        if (recordException.getRecord() != null) {
                            recordException.getRecord().receiverOffset().acknowledge();
                        }
                        repeatConsumer.set(true);
                    }
                })
                // 新增全局异常捕获,将异常终止转为正常结束,触发repeat
                .onErrorResume(throwable -> {
                    log.debug("Consumer stream error, will restart consumer: {}", throwable.getMessage());
                    repeatConsumer.set(true);
                    return Flux.empty();
                })
                .repeat(repeatConsumer::get);
    });
}

测试用例调整

调整Mock逻辑,验证重复订阅行为:

@Test
public void readWriteCreatesNewConsumerWhenCurrentConsumerStops() {
    AtomicInteger subscribeCount = new AtomicInteger(0);
    // 用thenAnswer模拟每次调用receiveAutoAck的返回值
    Mockito.when(reactiveKafkaConsumerTemplate.receiveAutoAck())
            .thenAnswer(invocation -> {
                int currentTimes = subscribeCount.getAndIncrement();
                if (currentTimes < 5) {
                    // 前5次订阅返回全局错误
                    return Flux.error(new RuntimeException("Kafka down"));
                } else {
                    // 第6次订阅返回正常记录
                    return Flux.just(createConsumerRecord(validMessage));
                }
            });

    Flux<String> actual = service.readWrite();

    StepVerifier.create(actual)
            // 期望收到正常返回的消息
            .expectNext(validMessage)
            .verifyComplete();

    // 验证订阅次数符合预期,确认repeat生效
    Assertions.assertEquals(6, subscribeCount.get());
}

核心原理

  • onErrorContinue是旁路错误处理机制,仅作用于单个元素处理过程中抛出的、和具体元素绑定的异常,上游源直接发出的全局Error信号不会触发该逻辑,会直接终止流。
  • repeat操作符的触发前提是上游流正常发出complete信号,因此需要用onErrorResume捕获全局异常,返回空Flux让流正常结束,同时设置重启标记。
  • Flux.defer的作用是每次重新订阅时都重新生成内部状态变量,避免上一次运行的状态残留导致逻辑错误。

内容的提问来源于stack exchange,提问作者perplexedDev

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 04:18:01