如何在Reactor Kafka消费时忽略不完整JSON消息?
解决方案
在Reactor流式处理中,要跳过特定异常场景的元素,不能用map(它必须返回一个值),而是要用flatMap——它支持返回Mono.empty()来过滤当前元素,同时保留正常流程和异常抛出的逻辑。
修改后的代码如下:
Flux<Person> consume() { return kafkaReceiver.receive() .flatMap(oneRecordWithJson -> { try { // 正常流程:转换JSON为Person对象,传递到下游 Person person = objectMapper.readValue(oneRecordWithJson.value(), Person.class); return Mono.just(person); } catch (JsonEOFException e) { // JSON不完整场景:记录日志后返回空Mono,直接跳过该消息 LOGGER.error("输入JSON不完整,忽略消息: {}", oneRecordWithJson.value(), e); return Mono.empty(); } catch (JsonProcessingException e) { // 其他JSON格式错误:抛出运行时异常,终止消费者 LOGGER.error("JSON格式错误,抛出异常终止消费: {}", oneRecordWithJson.value(), e); return Mono.error(new RuntimeException(e)); } }) .map(this::doSomething); }
关键逻辑说明:
flatMap替代map:flatMap允许返回Mono.empty(),这个空信号会被Reactor自动过滤,不会进入下游的doSomething处理;而map必须返回非空值,这就是原代码无法跳过元素的核心原因。- JSON不完整场景:捕获
JsonEOFException后返回Mono.empty(),实现“忽略消息”的需求,无需返回null或默认对象。 - 其他JSON错误场景:捕获
JsonProcessingException(除JsonEOFException外)后,用Mono.error()抛出运行时异常,触发消费者终止,符合需求。
补充说明(若需手动提交偏移量):
如果使用manual Ack模式,可在跳过消息时手动提交该记录的偏移量,避免重复消费:
catch (JsonEOFException e) { LOGGER.error("输入JSON不完整,忽略消息: {}", oneRecordWithJson.value(), e); oneRecordWithJson.receiveOffset().acknowledge(); return Mono.empty(); }
内容的提问来源于stack exchange,提问作者PatPanda
相关产品推荐
相关产品推荐

