如何基于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
相关产品推荐
相关产品推荐

