Nos artigos anteriores, apresentamos os conceitos fundamentais do RabbitMq e a integração com o SpringBoot. Agora, vamos explorar como utilizar o rabbitmq em um ambiente SpringBoot.
Este artigo focuses no envio de mensagens, incluindo:
- Uso básico do
RabbitTemplatepara envio de mensagens - Personalização de propriedades básicas da mensagem
- Personalização do conversor de mensagens
AbstractMessageConverter - Caso de falha ao enviar mensagens do tipo Object
I. Uso Básico
1. Configuração
Utilizaremos o SpringBoot 2.2.1.RELEASE com rabbitmq 3.7.5 para criar e testar o projeto.
Dependência no pom.xml:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
Configuração no arquivo application.yml:
spring:
rabbitmq:
virtual-host: /
username: admin
password: admin
port: 5672
host: 127.0.0.1
2. Classe de Configuração
Com base no conhecimento anterior sobre rabbitmq, sabemos que a lógica principal do remetente é "enviar a mensagem para o exchange e, em seguida, distribuí-la para a fila correspondente conforme diferentes estratégias".
Neste artigo, discutiremos principalmente o envio de mensagens. Para demonstrar os exemplos, definiremos um exchange do tipo topic e vincularemos uma fila.
public class ConstantesMensageria {
public static final String exchange = "topic.exchange";
public static final String routingKey = "chave.roteamento";
public final static String queue = "fila.topico";
}
@Configuration
public class ConfiguracaoMensageria {
@Bean
public TopicExchange exchangeTopico() {
return new TopicExchange(ConstantesMensageria.exchange);
}
@Bean
public Queue filaMensagens() {
return new Queue(ConstantesMensageria.queue, true);
}
@Bean
public Binding vinculacao(Queue filaMensagens, TopicExchange exchangeTopico) {
return BindingBuilder.bind(filaMensagens).to(exchangeTopico)
.with(ConstantesMensageria.routingKey);
}
@Bean
public RabbitTemplate modeloRabbit(ConnectionFactory factory) {
return new RabbitTemplate(factory);
}
}
3. Envio de Mensagens
O envio de mensagens é feito principalmente através do método RabbitTemplate#convertAndSend. Em geral, podemos usá-lo diretamente.
@Service
public class ServicoPublicador {
@Autowired
private RabbitTemplate modeloRabbit;
/**
* Uso comum - publicar mensagem
*
* @param dados Conteúdo da mensagem
* @return Mensagem publicada
*/
private String enviarMensagem(String dados) {
String mensagem = "Mensagem persistente = " + dados;
System.out.println("Publicando: " + mensagem);
modeloRabbit.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
mensagem
);
return mensagem;
}
}
O ponto principal está na linha: modeloRabbit.convertAndSend(ConstantesMensageria.exchange, ConstantesMensageria.routingKey, mensagem);
- Significa enviar a mensagem para o exchange especificado com a chave de roteamento definida
Observação
Pela abordagem acima, as mensagens enviadas são persistentes por padrão. Quando mensagens persistentes são distribuídas para filas persistentes, há operação de persistência em disco.
Em alguns cenários, não temos requisitos tão rigorosos para integridade dos dados, mas nos preocupamos mais com o desempenho do mq, e podemos aceitar a perda de alguns dados. Nesse caso, podemos personalizar as propriedades da mensagem (como definir a mensagem como não persistente).
Abaixo, duas abordagens são apresentadas, sendo a segunda recomendada.
/**
* Publicar mensagem não persistente. Quando publicada em fila persistente
* e o mq for reiniciado, esta mensagem será perdida.
*
* @param dados Dados da mensagem
* @return Mensagem publicada
*/
private String enviarMensagemNaoPersistente(String dados) {
MessageProperties propriedades = new MessageProperties();
propriedades.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
Message mensagem = modeloRabbit.getMessageConverter()
.toMessage("Mensagem_nao_persistente = " + dados, propriedades);
modeloRabbit.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
mensagem
);
System.out.println("Publicando: " + mensagem);
return mensagem.toString();
}
private String enviarMensagemPersonalizada(String dados) {
String msg = "Mensagem_customizada = " + dados;
modeloRabbit.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
msg,
new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
message.getMessageProperties().setHeader("cabecalho", "valor_teste");
message.getMessageProperties()
.setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT);
return message;
}
}
);
return msg;
}
Atenção
- No desenvolvimento real do projeto, recomenda-se usar
MessagePostProcessorpara personalizar propriedades da mensagem - Não é recomendado criar um objeto
MessagePostProcessorem cada envio de mensagem. Defina um objeto genérico que possa ser reutilizado
4. Caso de Erro ao Enviar Objeto Não Serializável
Analisando a interface rabbitTemplate#convertAndSend, sabemos que a mensagem enviada pode ser do tipo Object. Isso significa que qualquer objeto pode ser enviado para o mq?
Abaixo está um caso de teste:
private String enviarObjeto(String dados) {
EntidadeNaoSerializavel entidade = new EntidadeNaoSerializavel(18, dados);
System.out.println("Publicando: " + entidade);
modeloRabbit.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
entidade
);
return entidade.toString();
}
@Data
public static class EntidadeNaoSerializavel {
private Integer idade;
private String nome;
public EntidadeNaoSerializavel(int idade, String nome) {
this.idade = idade;
this.nome = nome;
}
}
Quando chamamos o método acima, não funcionará como esperado. Em vez disso, lançará uma exceção de tipo de parâmetro.
Por que isso acontece? Pela análise do stack trace, sabemos que o RabbitTemplate usa SimpleMessageConverter para implementar a lógica de encapsulamento da Message. O código principal é:
// Código de org.springframework.amqp.support.converter.SimpleMessageConverter#createMessage
@Override
protected Message createMessage(Object object, MessageProperties messageProperties)
throws MessageConversionException {
byte[] bytes = null;
if (object instanceof byte[]) {
bytes = (byte[]) object;
messageProperties.setContentType(MessageProperties.CONTENT_TYPE_BYTES);
}
else if (object instanceof String) {
try {
bytes = ((String) object).getBytes(this.defaultCharset);
}
catch (UnsupportedEncodingException e) {
throw new MessageConversionException(
"Falha ao converter para conteúdo Message", e);
}
messageProperties.setContentType(MessageProperties.CONTENT_TYPE_TEXT_PLAIN);
messageProperties.setContentEncoding(this.defaultCharset);
}
else if (object instanceof Serializable) {
try {
bytes = SerializationUtils.serialize(object);
}
catch (IllegalArgumentException e) {
throw new MessageConversionException(
"Falha ao converter para conteúdo Message serializado", e);
}
messageProperties.setContentType(MessageProperties.CONTENT_TYPE_SERIALIZED_OBJECT);
}
if (bytes != null) {
messageProperties.setContentLength(bytes.length);
return new Message(bytes, messageProperties);
}
throw new IllegalArgumentException(getClass().getSimpleName()
+ " suporta apenas String, byte[] e payloads Serializable, recebido: "
+ object.getClass().getName());
}
A lógica acima deixa claro que somente são aceitos arrays de bytes, strings e objetos sreializáveis (aqui使用的是 a serialização JDK para implementar a conversão entre objeto e array de bytes).
- Portanto, passar um objeto não serializável resultará em exceção.
Naturalmente, nos perguntamos se existe outro MessageConverter para suportar qualquer tipo de objeto de forma amigável.
5. Personalizando MessageConverter
A seguir, queremos resolver o problema acima personalizando um MessageConverter com serialização JSON.
Uma implementação simples (usando FastJson para serialização/desserialização):
public static class ConversorJsonPersonalizado extends AbstractMessageConverter {
@Override
protected Message createMessage(Object object, MessageProperties properties) {
properties.setContentType("application/json");
return new Message(JSON.toJSONBytes(object), properties);
}
@Override
public Object fromMessage(Message message) throws MessageConversionException {
return JSON.parse(message.getBody());
}
}
Redefinindo um rabbitTemplate e definindo seu conversor de mensagens como o ConversorJsonPersonalizado:
@Bean
public RabbitTemplate modeloRabbitJson(ConnectionFactory factory) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(factory);
rabbitTemplate.setMessageConverter(new ConversorJsonPersonalizado());
return rabbitTemplate;
}
Agora testando novamente:
@Service
public class PublicadorJson {
@Autowired
private RabbitTemplate modeloRabbitJson;
private String publicarMapa(String dados) {
Map<String, Object> msg = new HashMap<>(8);
msg.put("mensagem", dados);
msg.put("tipo", "json");
msg.put("versao", 123);
System.out.println("Publicando: " + msg);
modeloRabbitJson.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
msg
);
return msg.toString();
}
private String publicarObjeto(String dados) {
EntidadeNaoSerializavel entidade =
new EntidadeNaoSerializavel(18, "JSON_PROPRIO" + dados);
System.out.println("Publicando: " + entidade);
modeloRabbitJson.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
entidade
);
return entidade.toString();
}
}
A mensagem recebida no mq será exibida corretamente.
6. Jackson2JsonMessageConverter
Embora tenhamos alcançado a conversão de formato JSON acima, a implementação é rudimentar. E como essa funcionalidade básica e genérica segue o padrão do ecossistema Spring, com certeza existe uma solução pronta. Estamos falando do Jackson2JsonMessageConverter.
Nossa forma de uso pode ser a seguinte:
//Definir RabbitTemplate
@Bean
public RabbitTemplate modeloRabbitJackson(ConnectionFactory factory) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(factory);
rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter());
return rabbitTemplate;
}
// Código de teste
@Autowired
private RabbitTemplate modeloRabbitJackson;
private String publicarComJackson(String dados) {
Map<String, Object> msg = new HashMap<>(8);
msg.put("mensagem", dados);
msg.put("tipo", "jackson");
msg.put("versao", 456);
System.out.println("Publicando: " + msg);
modeloRabbitJackson.convertAndSend(
ConstantesMensageria.exchange,
ConstantesMensageria.routingKey,
msg
);
return msg.toString();
}
Abaixo está o conteúdo da mensagem após serialização com Jackson. Diferente da nossa implementação personalizada, há campos adicionais como headers e content_encoding.
7. Resumo
Os principais conhecimentos deste artigo são:
- Usar
RabbitTemplate#convertAndSendpara分发mensagens - Usar
MessagePostProcessorpara personalizar propriedades da mensagem (observe que por padrão as mensagens são persistentes) - O encapsulador padrão de mensagens é
SimpleMessageConverter, que suporta apenas分发arrays de bytes, strings e objetos serializáveis; chamadas de método que não atendam a essas três condições lançarão exceção - Podemos implementar a interface
MessageConverterpara definir nosso próprio encapsulador de mensagens e resolver o problema acima
Nos artigos sobre RabbitMq, mencionamos que para garantir que a mensagem seja recebida corretamente pelo broker, existem dois mecanismos: confirmação de mensagem e transação. Como o produtor de mensagens deve usar esses dois métodos?
Devido ao limite de篇幅, o próximo artigo apresentará o uso do envio de mensagens sob o mecanismo de confirmação/transação.
II. Outros
0. Artigos Relacionados & Código Fonte
Artigos da série
- 【MQ Series】springboot + rabbitmq First Experience
- 【MQ Series】RabbitMq Core Knowledge Summary
Código fonte do projeto