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(); }
疑问解答
MultiEmitterProcessor是否适合作为热流使用?
MultiEmitterProcessor本身就是热流处理器,它创建后会立即开始接收事件,即使没有订阅者,事件也会按默认策略(保留最后1个)缓存。它的设计目标就是支持多订阅场景下的事件广播,完全适合作为热流使用。另外你代码里的toHotStream()调用是多余的,因为MultiEmitterProcessor的toMulti()返回的已经是热流。是否存在类似Reactor的ConnectableFlux可手动连接热流?
Smallrye Mutiny中对应的是ConnectableMulti,可以通过Multi#toConnectable()方法将普通Multi转换为可手动控制连接的热流。你可以先创建ConnectableMulti实例,在需要启动事件发送和订阅时调用connect()方法,以此实现手动控制热流的连接时机,避免提前消费事件。
内容的提问来源于stack exchange,提问作者Hantsy
相关产品推荐
相关产品推荐

