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

如何基于Quarkus Messaging实现IBM MQ消息转换异常处理?

基于Quarkus Messaging + SmallRye JMS实现IBM MQ消息转换与异常分流

1. 依赖配置

在pom.xml中添加必要依赖,包含Quarkus Reactive Messaging JMS连接器、IBM MQ客户端、JSON处理及可选的校验组件:

<dependencies>
    <!-- Quarkus Reactive Messaging JMS 连接器 -->
    <dependency>
        <groupId>io.quarkus</groupId>
        <artifactId>quarkus-smallrye-reactive-messaging-jms</artifactId>
    </dependency>
    <!-- IBM MQ 客户端 -->
    <dependency>
        <groupId>com.ibm.mq</groupId>
        <artifactId>com.ibm.mq.allclient</artifactId>
    </dependency>
    <!-- Jackson JSON 处理 -->
    <dependency>
        <groupId>io.quarkus</groupId>
        <artifactId>quarkus-jackson</artifactId>
    </dependency>
    <!-- 可选:Jakarta 校验组件(用于业务对象字段校验) -->
    <dependency>
        <groupId>io.quarkus</groupId>
        <artifactId>quarkus-hibernate-validator</artifactId>
    </dependency>
    <!-- SLF4J 日志 -->
    <dependency>
        <groupId>io.quarkus</groupId>
        <artifactId>quarkus-logging-slf4j</artifactId>
    </dependency>
</dependencies>

2. 配置文件

在application.properties中配置IBM MQ连接信息及Reactive Messaging通道映射:

# IBM MQ JMS 连接配置
quarkus.jms.url=tcp://你的MQ主机:1414
quarkus.jms.username=MQ用户名
quarkus.jms.password=MQ密码
quarkus.jms.connection-factory-name=com.ibm.mq.jms.MQConnectionFactory

# Reactive Messaging 通道配置
# 输入通道:监听queue_A
mp.messaging.incoming.queue_A.connector=smallrye-jms
mp.messaging.incoming.queue_A.destination=queue_A
mp.messaging.incoming.queue_A.destination-type=queue

# 输出通道:消息转换成功后发送到queue_B
mp.messaging.outgoing.queue_B.connector=smallrye-jms
mp.messaging.outgoing.queue_B.destination=queue_B
mp.messaging.outgoing.queue_B.destination-type=queue

# 输出通道:转换失败时发送到queue_A_error
mp.messaging.outgoing.queue_A_error.connector=smallrye-jms
mp.messaging.outgoing.queue_A_error.destination=queue_A_error
mp.messaging.outgoing.queue_A_error.destination-type=queue

3. 消息转换器实现

封装JSON与Java对象的转换逻辑,同时处理转换及校验异常:

import com.fasterxml.jackson.databind.ObjectMapper;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.inject.Inject;
import jakarta.validation.ConstraintViolation;
import jakarta.validation.Validator;
import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;

@ApplicationScoped
public class MessageConverter {

    private final ObjectMapper objectMapper = new ObjectMapper();

    // 可选:注入校验器用于业务对象字段校验
    @Inject
    Validator validator;

    /**
     * 将字符串消息转换为目标Java对象,返回转换结果(含数据或错误列表)
     */
    public ConversionResult convertToObject(String payload, Class<?> targetClass) {
        try {
            Object result = objectMapper.readValue(payload, targetClass);
            // 执行业务对象字段校验(可选)
            Set<ConstraintViolation<Object>> violations = validator.validate(result);
            if (!violations.isEmpty()) {
                List<String> errors = violations.stream()
                        .map(v -> v.getPropertyPath() + ": " + v.getMessage())
                        .collect(Collectors.toList());
                return new ConversionResult(null, errors);
            }
            return new ConversionResult(result, null);
        } catch (Exception e) {
            // 捕获JSON解析异常,返回错误信息
            return new ConversionResult(null, List.of("JSON解析失败: " + e.getMessage()));
        }
    }

    /**
     * 将Java对象序列化为JSON字符串
     */
    public String convertToJson(Object obj) throws Exception {
        return objectMapper.writeValueAsString(obj);
    }

    // 封装转换结果的内部类
    public static class ConversionResult {
        public final Object data;
        public final List<String> errors;

        public ConversionResult(Object data, List<String> errors) {
            this.data = data;
            this.errors = errors;
        }
    }
}

4. 消息路由与异常分流

实现消息消费者,监听queue_A,根据转换结果分流到对应队列:

import jakarta.enterprise.context.ApplicationScoped;
import jakarta.inject.Inject;
import org.eclipse.microprofile.reactive.messaging.Emitter;
import org.eclipse.microprofile.reactive.messaging.Incoming;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.Map;

@ApplicationScoped
public class MessageRouter {

    private static final Logger LOGGER = LoggerFactory.getLogger(MessageRouter.class);

    @Inject
    MessageConverter converter;

    // 注入发送到queue_B的发射器
    @Inject
    @Outgoing("queue_B")
    Emitter<String> queueBEmitter;

    // 注入发送到错误队列的发射器
    @Inject
    @Outgoing("queue_A_error")
    Emitter<String> errorQueueEmitter;

    /**
     * 处理来自queue_A的消息
     */
    @Incoming("queue_A")
    public void processMessage(String payload) {
        // 替换为你的业务对象类
        MessageConverter.ConversionResult result = converter.convertToObject(payload, YourBusinessObject.class);

        if (result.errors != null && !result.errors.isEmpty()) {
            // 转换失败,构造错误消息并发送到queue_A_error
            Map<String, Object> errorMsg = new HashMap<>();
            errorMsg.put("originalPayload", payload);
            errorMsg.put("errors", result.errors);
            try {
                String errorJson = converter.convertToJson(errorMsg);
                errorQueueEmitter.send(errorJson);
                LOGGER.info("已发送错误消息到queue_A_error");
            } catch (Exception e) {
                LOGGER.error("发送错误消息失败", e);
            }
        } else {
            // 转换成功,序列化对象后发送到queue_B
            try {
                String targetJson = converter.convertToJson(result.data);
                queueBEmitter.send(targetJson);
                LOGGER.info("已处理消息并发送到queue_B");
            } catch (Exception e) {
                LOGGER.error("发送处理后消息到queue_B失败", e);
            }
        }
    }
}

5. 业务对象定义

创建对应的业务对象类,可添加校验注解(可选):

import jakarta.validation.constraints.NotEmpty;

public class YourBusinessObject {
    @NotEmpty(message = "ID不能为空")
    private String id;

    @NotEmpty(message = "名称不能为空")
    private String name;

    // 省略getter、setter方法
}

关键注意点

  • 消息确认:默认配置下,Quarkus会在消息处理完成后自动确认MQ消息;若需手动控制,可添加mp.messaging.incoming.queue_A.auto-acknowledge=false,并结合@Acknowledgment注解实现手动确认。
  • 异常重试:若发送错误消息或目标消息失败,可通过配置通道的重试策略(如mp.messaging.outgoing.queue_B.retry.max-attempts=3)实现自动重试。
  • 日志记录:务必在异常分支添加日志,便于排查问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 15:27:37