如何使用Smallrye Messaging与响应式HTTP客户端实现非阻塞处理
问题解答
你的核心思路方向是对的,但给出的实现代码存在错误,同时还有更符合Reactive Messaging规范的简化写法,不需要手动写订阅逻辑。
原阻塞异常的触发逻辑
Reactive Messaging框架默认在非阻塞的事件循环线程上执行消费方法,事件循环线程禁止执行任何阻塞操作。你最初的同步客户端调用直接返回String类型,属于阻塞调用,因此触发阻塞线程异常。添加@Blocking注解后,框架会将该方法调度到专属的阻塞工作线程池执行,因此异常消失。
客户端改为返回Uni<String>后的正确实现
当客户端方法改为返回响应式类型Uni<String>后,你不需要再添加@Blocking注解,也不用强制切换手动ack模式,框架原生支持响应式返回值的自动ack/nack,分两种场景处理:
- 无特殊消息元数据需求的场景:直接入参传消息payload(String类型),方法返回
Uni<Void>,框架会自动订阅返回的Uni实例:Uni成功完成时自动ack消息,Uni抛出异常时自动nack消息,无需手动处理确认逻辑,代码如下:
@Incoming("my-topic") Uni<Void> process(String someEvent) { return assuminglyRealFastReactiveClientCall() .onItem().invoke(resp -> System.out.println("Response: " + resp)) .replaceWithVoid(); }
- 需要获取消息元数据做自定义处理的场景:才需要将入参改为
Message<String>类型。你之前写的手动订阅代码存在两个问题:- 变量名不匹配:入参命名为
someEvent,代码里调用ack/nack时用的是event,会直接编译失败 - 方法内部手动调用
subscribe()属于不规范写法,会绕过框架的背压控制、上下文传递、异常兜底逻辑,容易出现消息重复消费、异常丢失的问题。正确的手动处理写法是将ack/nack逻辑锚定到响应式流中,把构建好的Uni返回给框架接管订阅:
- 变量名不匹配:入参命名为
@Incoming("my-topic") Uni<Void> process(Message<String> someEvent) { return assuminglyRealFastReactiveClientCall() .onItem().invoke(resp -> System.out.println("Response: " + resp)) .onItem().transformToUni(resp -> someEvent.ack()) .onFailure().recoverWithUni(err -> someEvent.nack(err)); }
注意:不要在
@Incoming修饰的消费方法内部手动执行订阅操作,框架本身就是响应式流的最终订阅者,只需要将构建完成的响应式流(Uni/Multi)返回给框架即可,框架会统一处理线程调度、异常兜底、消息确认全流程。
内容的提问来源于stack exchange,提问作者user432024
相关产品推荐
相关产品推荐

