基于Quarkus的Tibco EMS消息消费:如何用Mutiny实现响应式改造?
问题
我正在开发一个Quarkus应用,需要从Tibco EMS消费消息并进行处理。以下是我当前的实现代码片段:
Session session = connection.createSession(); Destination destination = session.createQueue(configurationProperties.queueFlight()); MessageConsumer consumer = session.createConsumer(destination); consumer.setMessageListener(message -> { FlightData flightData = flightDataConverter.toFlightData(message); // TODO: send to kafka });
这段代码是Quarkus应用启动时调用的方法的一部分,可成功接收Tibco EMS的消息。但我希望通过Mutiny实现响应式方案优化流程,现求教是否可行及具体改造方法,欢迎提供建议或示例。
解决方案
完全可行,Quarkus与Mutiny深度集成,结合Tibco EMS可以轻松实现响应式消息消费流程,以下是具体的改造思路和示例代码:
1. 手动包装Listener为Mutiny流
如果直接使用原生JMS API,可以将传统的MessageListener包装成Mutiny的Multi(代表0到N个元素的响应式流),从而利用Mutiny的操作符处理消息:
@ApplicationScoped public class TibcoEmsReactiveConsumer { private final ConfigurationProperties configProps; private final FlightDataConverter flightConverter; private final KafkaProducer reactiveKafkaProducer; // 假设使用Quarkus响应式Kafka生产者 // 构造函数注入依赖 public TibcoEmsReactiveConsumer(ConfigurationProperties configProps, FlightDataConverter flightConverter, KafkaProducer reactiveKafkaProducer) { this.configProps = configProps; this.flightConverter = flightConverter; this.reactiveKafkaProducer = reactiveKafkaProducer; } @PostConstruct void startReactiveConsumer() { try { Session session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE); Destination queue = session.createQueue(configProps.queueFlight()); MessageConsumer consumer = session.createConsumer(queue); // 将MessageListener转换为Mutiny Multi Multi<Message> messageStream = Multi.createFrom().emitter(emitter -> { consumer.setMessageListener(message -> { try { emitter.emit(message); message.acknowledge(); // 手动确认消息,根据Session模式调整 } catch (JMSException e) { emitter.fail(e); } }); // 注册终止回调,清理资源 emitter.onTermination(() -> { try { consumer.close(); session.close(); } catch (JMSException e) { // 记录关闭异常日志 } }); }); // 处理消息流:转换 -> 发送Kafka -> 错误处理 messageStream .onItem().transform(this::convertToFlightData) .onItem().transformToUni(flightData -> reactiveKafkaProducer.send("flight-topic", flightData)) .onFailure().retry().atMost(3) // 失败重试3次 .onFailure().invoke(failure -> { /* 记录最终失败日志 */ }) .subscribe().with( result -> {}, // 处理Kafka发送成功的结果 failure -> {} // 处理无法恢复的错误 ); } catch (JMSException e) { throw new RuntimeException("Failed to initialize Tibco EMS consumer", e); } } private FlightData convertToFlightData(Message message) { try { return flightConverter.toFlightData(message); } catch (Exception e) { throw new RuntimeException("Failed to convert message to FlightData", e); } } }
2. 使用Quarkus JMS扩展简化实现
如果你的项目引入了Quarkus的JMS扩展(适配Tibco EMS),可以直接使用扩展提供的Mutiny风格API,代码会更简洁:
@ApplicationScoped public class QuarkusTibcoReactiveConsumer { @Inject @JmsConnection("tibco-ems") // 对应application.properties中的JMS连接配置 MutinyJmsConnectionFactory connectionFactory; private final FlightDataConverter flightConverter; private final KafkaProducer reactiveKafkaProducer; public QuarkusTibcoReactiveConsumer(FlightDataConverter flightConverter, KafkaProducer reactiveKafkaProducer) { this.flightConverter = flightConverter; this.reactiveKafkaProducer = reactiveKafkaProducer; } @PostConstruct void initReactiveConsumer() { connectionFactory.createSession() .flatMap(session -> session.createConsumer(configProps.queueFlight())) .subscribe().with(consumer -> { // 订阅消息流并处理 consumer.messages() .onItem().transform(this::convertToFlightData) .onItem().transformToUni(this::sendToKafka) .onFailure().retry().withBackOff(Duration.ofMillis(500)) // 指数退避重试 .onFailure().invoke(failure -> { /* 记录错误 */ }) .subscribe().with( ignored -> {}, failure -> {} ); }); } private Uni<Void> sendToKafka(FlightData flightData) { return reactiveKafkaProducer.send("flight-topic", flightData) .onItem().ignore().andContinueWithNull(); } private FlightData convertToFlightData(Message message) { try { return flightConverter.toFlightData(message); } catch (Exception e) { throw new RuntimeException("Message conversion failed", e); } } }
3. 核心注意点
- 消息确认机制:根据Session的事务模式(事务性/非事务性)调整消息确认逻辑,确保消息不丢失、不重复消费。
- 资源管理:通过Mutiny的
onTermination或Quarkus的@PreDestroy回调清理JMS资源,避免资源泄漏。 - 错误处理:利用Mutiny的
onFailure系列操作符实现重试、降级或错误日志,提升系统容错能力。 - 背压支持:Mutiny的
Multi天然支持背压,当下游处理能力不足时会自动调节消息接收速率,防止内存溢出。
内容的提问来源于stack exchange,提问作者andrean
相关产品推荐
相关产品推荐

