Quarkus响应式消息用Message签名处理器时ignore失败策略报循环引用错误
根本原因
当处理器方法显式使用Message<T>作为入参和返回值类型时,SmallRye Reactive Messaging会默认你需要完全自主控制消息的生命周期(包含确认、否定确认、异常处理),此时框架不会自动拦截方法抛出的异常,异常会直接向上传播导致流终止,你配置的ignore失败策略不会自动触发。
而使用普通payload类型(比如示例中的String)作为入参、返回值时,框架会自动封装消息的生命周期管理逻辑,方法抛出异常时会自动匹配配置的失败策略,记录日志后跳过异常消息,消费流可以正常继续。
解决方案
方案1:手动处理异常完成消息生命周期
调整处理器方法,用try-catch包裹业务逻辑,出现异常时主动调用消息的nack()方法,不需要向下游传递消息时直接返回null即可:
@Incoming("movies") @Outgoing("next") public Message<String> consume(Message<String> movie) { LOGGER.infof("Receiving movie %s", movie); try { if (movie.getPayload().contains("'")) { throw new IllegalArgumentException("I don't like movie with ' in their title: " + movie); } if (movie.getPayload().contains(",")) { throw new IllegalArgumentException("I don't like movie with , in their title: " + movie); } return movie.withPayload(movie.getPayload()); } catch (Exception e) { // 手动否定确认异常消息 movie.nack(e); // 不向下游传递当前异常消息 return null; } }
方案2:使用响应式返回类型让框架接管异常处理
如果不想手动编写异常处理逻辑,可以将返回值改为Uni<Message<T>>类型,此时框架会自动感知方法执行的异常,触发你配置的ignore失败策略,不需要手动管理ack/nack逻辑:
@Incoming("movies") @Outgoing("next") public Uni<Message<String>> consume(Message<String> movie) { LOGGER.infof("Receiving movie %s", movie); if (movie.getPayload().contains("'")) { throw new IllegalArgumentException("I don't like movie with ' in their title: " + movie); } if (movie.getPayload().contains(",")) { throw new IllegalArgumentException("I don't like movie with , in their title: " + movie); } return Uni.createFrom().item(movie.withPayload(movie.getPayload())); }
请确保application.properties中已经配置了正确的失败策略:
mp.messaging.incoming.movies.failure-strategy=ignore
内容的提问来源于stack exchange,提问作者Maxime
相关产品推荐
相关产品推荐

