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

Quarkus中JMS消息转SSE端点实现问题求助

关于Quarkus中Smallrye Mutiny热流的疑问解答

在Spring项目中使用Sinks向SSE端点发送事件运行正常,但在Quarkus中尝试用Smallrye Mutiny的MultiEmitterProcessor实现相同功能时遇到问题,相关代码如下:

JMS消息接收代码

MultiEmitterProcessor<Message> emitterProcessor = MultiEmitterProcessor.create();

void receive() {
    var consumer = jmsContext.createConsumer(helloQueue);
    consumer.setMessageListener(
            msg -> {
                try {
                    var received = jsonb.fromJson(msg.getBody(String.class), Message.class);
                    LOGGER.log(Level.INFO, "consuming message: {0}", received);
                    emitterProcessor.emit(received);
                } catch (JMSException e) {
                    throw new RuntimeException(e);
                }
            }
    );
}

Resource类代码

@GET
@Produces(MediaType.SERVER_SENT_EVENTS)
@RestStreamElementType(MediaType.APPLICATION_JSON)
public Multi<Message> stream() {
    return handler.emitterProcessor.toMulti().toHotStream();
}

疑问解答

  1. MultiEmitterProcessor是否适合作为热流使用?
    MultiEmitterProcessor本身就是热流处理器,它创建后会立即开始接收事件,即使没有订阅者,事件也会按默认策略(保留最后1个)缓存。它的设计目标就是支持多订阅场景下的事件广播,完全适合作为热流使用。另外你代码里的toHotStream()调用是多余的,因为MultiEmitterProcessor的toMulti()返回的已经是热流。

  2. 是否存在类似Reactor的ConnectableFlux可手动连接热流?
    Smallrye Mutiny中对应的是ConnectableMulti,可以通过Multi#toConnectable()方法将普通Multi转换为可手动控制连接的热流。你可以先创建ConnectableMulti实例,在需要启动事件发送和订阅时调用connect()方法,以此实现手动控制热流的连接时机,避免提前消费事件。


内容的提问来源于stack exchange,提问作者Hantsy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 06:16:15