Spring @JmsListener事务场景下如何阻止特定消息重投
针对事务会话下@JmsListener参数转换阶段抛出异常导致非预期消息重投的问题,你之前尝试的全局ErrorHandler方案不生效的核心原因是:容器级别的全局ErrorHandler触发时机在事务回滚操作完成之后,此时就算不再向外抛出异常,回滚和重投逻辑已经执行完毕,自然无法达到预期效果。
除了你提到的手动接收TextMessage自行做序列化的方案外,还有以下几种侵入性更低的实现方案:
核心思路是把消息转换阶段的异常捕获逻辑收口在消息转换层,不改动原有监听器的方法签名,实现逻辑和业务逻辑完全解耦。
你当前应该是使用MappingJackson2MessageConverter做JSON消息和实体类的转换,只需要继承该类重写fromMessage方法,捕获所有JSON解析、类型映射相关的异常,做统一处理即可:
@Slf4j public class SafeJacksonMessageConverter extends MappingJackson2MessageConverter { // 定义非法消息标记常量 public static final Object INVALID_MESSAGE = new Object(); @Override public Object fromMessage(Message message) throws JMSException, MessageConversionException { try { return super.fromMessage(message); } catch (MessageConversionException | JsonProcessingException e) { // 提取原始消息内容留痕 String rawPayload = message instanceof TextMessage ? ((TextMessage) message).getText() : message.toString(); log.error("消息格式不合法,不再重试投递,原始消息:{}", rawPayload, e); // 可选操作:将非法消息转发到专属死信队列,后续集中处理 // jmsTemplate.convertAndSend("DLQ.my-destination", rawPayload); return INVALID_MESSAGE; } } }
将这个自定义转换器注册到你使用的my-container-factory中,原有监听器只需要增加一行非法消息判断即可,不需要修改方法签名:
@JmsListener(destination = "my-destination", containerFactory = "my-container-factory") void receive(MyClass myObject) { if (myObject == SafeJacksonMessageConverter.INVALID_MESSAGE) { // 非法消息直接返回,事务正常提交,不会触发重投 return; } // 原有业务逻辑 // ... }
该方案的优点是逻辑统一收口,不会对原有业务代码造成大面积侵入,同时可以灵活定制非法消息的处理逻辑(日志记录、落库、转死信等),适合绝大多数场景。
如果你的项目中所有使用该容器工厂的监听器,都不需要对消息转换异常做重试,可以直接在DefaultJmsListenerContainerFactory中配置不触发回滚的异常类型,搭配全局错误处理器记录日志即可:
@Bean public DefaultJmsListenerContainerFactory my-container-factory(ConnectionFactory connectionFactory) { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setSessionTransacted(true); // 配置消息转换异常不触发事务回滚 factory.setTransactionNoRollbackFor(MessageConversionException.class); // 全局异常记录 factory.setErrorHandler(t -> log.error("JMS监听器消费异常", t)); // 注册原有消息转换器 factory.setMessageConverter(new MappingJackson2MessageConverter()); return factory; }
该方案配置最简单,但是作用范围是整个容器工厂下的所有监听器,如果个别监听器需要对转换异常做特殊重试逻辑,会无法兼容,适合规则统一的简单项目使用。
@JmsListener注解支持单独为当前监听器指定专属的errorHandler,该处理器的触发时机早于事务回滚操作,处理完成后如果不向外抛出异常,事务会正常提交,不会触发消息重投:
首先自定义方法级错误处理器:
@Slf4j public class ListenerErrorHandler { public void handleError(Message message, ListenerExecutionFailedException e) throws JMSException { Throwable cause = e.getCause(); // 只处理消息转换类异常,其他异常继续抛出触发回滚 if (cause instanceof MessageConversionException || cause.getCause() instanceof JsonProcessingException) { String payload = message instanceof TextMessage ? ((TextMessage) message).getText() : message.toString(); log.error("消息格式非法,不做重试,原始消息:{}", payload, e); return; } // 其他异常继续抛出,走原有回滚重试逻辑 throw e; } }
在注解中指定该处理器即可:
@JmsListener( destination = "my-destination", containerFactory = "my-container-factory", errorHandler = "listenerErrorHandler" ) void receive(MyClass myObject) { // 原有业务逻辑无需改动 // ... }
该方案的粒度最细,可以针对单个监听器定制异常处理规则,不需要修改全局配置,但是需要注意区分异常类型,避免把需要重试的业务异常也一并吞掉。
注意:所有方案处理非法消息时,一定要做好原始消息和异常信息的留痕,有条件的建议统一转发到死信队列做后续排查,不要直接静默丢弃消息,避免出现消息丢失无据可查的问题。
内容的提问来源于stack exchange,提问作者Yonas

