WebFlux Reactor Kafka重试机制未触发processRecord方法问题排查
问题根因
你的重试机制不生效的核心原因是**retryWhen的挂载位置错误,重试作用域覆盖了整个分区消费流,而非单条消息的处理逻辑**。
当前代码中retryWhen直接拼接在分区GroupedFlux的处理链尾部,触发重试时会重新订阅整个分区的消费流。但groupBy生成的分区消费GroupedFlux不会重放已经推送过的消息:首次执行时offset=937的失败消息已经从流中发出,重试重订阅流时只会拉取后续新到达的消息,不会重新推送之前已经发过的失败消息,因此重试阶段不会再次调用processRecord处理这条失败消息,你看到的3次重试只是反复重订阅没有目标消息的空流,直到重试次数耗尽抛出异常。
从日志也能佐证这个问题:三次重试触发的过程中,始终没有再次打印Record received from offset...的消息接收日志,说明重试阶段根本没有拿到之前失败的那条消息。
修复方案
将重试逻辑的作用域收紧到单条消息维度,让重试仅针对单条消息的processRecord处理流程生效,避免重订阅整个分区流。
核心修改后的代码片段如下:
.flatMap(partitionFlux -> partitionFlux.subscribeOn(scheduler) // 单条消息维度处理,把重试逻辑包在单条消息的flatMap内部 .flatMap(r -> Mono.just(r) .doOnNext(record -> log.info("processEvent:: Record received from offset {} from topicPartition {} with message key {}", record.receiverOffset().topicPartition(), record.key(), record.offset())) .flatMap(this::processRecord) .doOnNext(receiverRecordInfo -> log.info("processEvent:: Record processed from offset {} from topicPartition {} with message key {}", receiverRecordInfo.receiverOffset().offset(), receiverRecordInfo.receiverOffset().topicPartition())) // 重试仅作用于当前单条消息的处理 .retryWhen(Retry.backoff(3, Duration.ofMillis(200)) .jitter(0d) .doAfterRetry(retrySignal -> log.error("Retried {} for record offset {}", retrySignal.totalRetries(), r.offset())) .onRetryExhaustedThrow((retryBackoffSpec, retrySignal) -> new KafkaRetryExhaustException()) ) ) .doOnError(KafkaRetryExhaustException.class, (msg) -> { log.error("exception... ", msg); }) .onErrorResume(throwable -> Mono.empty()) )
修改后每次触发重试,都会针对当前失败的消息重新执行processRecord方法,符合你的预期。需要注意不要让单条消息的处理异常冒泡到分区流层面,否则会导致整个分区的消费流中断,影响后续消息消费。
内容的提问来源于stack exchange,提问作者Rambabu Kokkiligadda
相关产品推荐
相关产品推荐

