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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 20:57:16