Spring AMQP:无法结合RabbitTemplate使用ContentTypeDelegatingMessageConverter
解决ContentTypeDelegatingMessageConverter与RabbitTemplate配合问题
核心问题在于RabbitTemplate.convertAndSend()的执行顺序:先调用MessageConverter完成消息转换(此时使用默认空MessageProperties,contentType为octet-stream),再执行MessagePostProcessor设置contentType,导致ContentTypeDelegatingMessageConverter无法根据目标contentType选择对应转换器。
以下是两种可行的解决方案:
方案1:手动调用MessageConverter构造消息后发送
直接通过MessageConverter使用指定的contentType构造Message对象,再调用rabbitTemplate.send()发送,确保转换时就能拿到正确的contentType:
@Service public class RabbitMQService { @Autowired private RabbitTemplate rabbitTemplate; @Autowired private MessageConverter messageConverter; public void sendMessage(final String exchange, final String routingKey, final Object payload, String contentType, boolean persistent) { // 提前构造带正确contentType的MessageProperties MessageProperties props = new MessageProperties(); props.setContentType(contentType); props.setDeliveryMode(persistent ? MessageDeliveryMode.PERSISTENT : MessageDeliveryMode.NON_PERSISTENT); // 用指定属性转换消息 Message message = messageConverter.toMessage(payload, props); rabbitTemplate.send(exchange, routingKey, message); } }
方案2:扩展RabbitTemplate,新增支持传入MessageProperties的重载方法
通过扩展RabbitTemplate,添加一个能直接传入MessageProperties的convertAndSend重载方法,让转换阶段就使用指定的属性:
1. 自定义RabbitTemplate子类
public class CustomRabbitTemplate extends RabbitTemplate { public CustomRabbitTemplate(ConnectionFactory connectionFactory) { super(connectionFactory); } public void convertAndSend(String exchange, String routingKey, Object payload, MessageProperties messageProperties, @Nullable MessagePostProcessor messagePostProcessor) throws AmqpException { // 使用传入的MessageProperties进行转换 Message messageToSend = getRequiredMessageConverter().toMessage(payload, messageProperties); if (messagePostProcessor != null) { messageToSend = messagePostProcessor.postProcessMessage(messageToSend, null, nullSafeExchange(exchange), nullSafeRoutingKey(routingKey)); } send(exchange, routingKey, messageToSend); } }
2. 替换配置中的RabbitTemplate实例
@Bean public RabbitTemplate rabbitTemplate(final ConnectionFactory connectionFactory, final Jackson2JsonMessageConverter jsonMessageConverter, final Jackson2XmlMessageConverter xmlMessageConverter) { // 使用自定义的CustomRabbitTemplate final var rabbitTemplate = new CustomRabbitTemplate(connectionFactory); ContentTypeDelegatingMessageConverter compositeConverter = new ContentTypeDelegatingMessageConverter(); compositeConverter.addDelegate(MessageProperties.CONTENT_TYPE_JSON, jsonMessageConverter); compositeConverter.addDelegate(MessageProperties.CONTENT_TYPE_XML, xmlMessageConverter); rabbitTemplate.setMessageConverter(compositeConverter); return rabbitTemplate; }
3. 在服务中调用新方法
@Service public class RabbitMQService { @Autowired private CustomRabbitTemplate rabbitTemplate; public void sendMessage(final String exchange, final String routingKey, final Object payload, String contentType, boolean persistent) { MessageProperties props = new MessageProperties(); props.setContentType(contentType); props.setDeliveryMode(persistent ? MessageDeliveryMode.PERSISTENT : MessageDeliveryMode.NON_PERSISTENT); // 调用自定义的convertAndSend方法 rabbitTemplate.convertAndSend(exchange, routingKey, payload, props, null); } }
内容的提问来源于stack exchange,提问作者M.Ricciuti
相关产品推荐
相关产品推荐

