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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 08:52:35