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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 21:58:30