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
相关产品推荐
相关产品推荐

