Conceitos e Uso do RabbitMQ

Protocolo AMQP explicado, site oficial: https://www.rabbitmq.com/tutorials/amqp-concepts

O RabbitMQ é um middleware de filas de mensagens que utiliza por padrão o protocolo AMQP 0-9-1. AMQP significa Advanced Message Queuing Protocol, ou seja, Protocolo Avançado de Filas de Mensagens. O Spring oferece suporte ao RabbitMQ através do módulo spring-amqp.

Os componentes fundamentais do protocolo são três: 1) Publicador (publisher), 2) Agente de Mensagens (messaging broker), 3) Consumidor (consumer).

Dentro do agente de mensagens existem três conceitos principais: 1) Fila (queue), 2) Exchange (intercâmbio), 3) Ligações (bindings).

O fluxo de comunicação segue este caminho: o publicador envia uma mansagem ao exchange dentro do agente, e este, com base nas regras de ligação, encaminha a mensagem para as filas associadas. Posteriormente, o consumidor recebe essas mensagens via push ou pull API.

I. Exchanges

Um Exchange é um dos conceitos mais importantes no sistema. Um agente pode possuir múltiplos exchanges, mas cada publicador só se conecta a um único, enviando suas mensagens exclusivamente para ele. O exchange não armazena as mensagens, apenas direciona-as às filas conforme as regras configuradas.

Existem quatro tipos de exchanges:

  1. Direct Exchange: Utiliza uma chave de roteamento única. A mensagem é entregue somente quando a chave de roteamento do publicador e do consumidor forem idênticas.
  2. Fanout Exchange: Não considera a chave de roteamento. Todas as filas vinculadas recebem a mensagem.
  3. Topic Exchange: Permite roteamento baseado em padrões com uso de asteriscos (* e #) para correspondência parcial.
  4. Headers Exchange: Roteia mensagens baseado nos cabeçalhos da mensagem em vez da chave de roteamento.

Além disso, os exchanges têm propriedades importantes como nome, durabilidade (persistência após reinicialização do broker), e exclusividade automática (remoção quando nenhuma fila estiver ligada).

II. Filas (Queues)

As filas são estruturas FIFO (First In, First Out) onde as mensagens são armazenadas. Para utilizá-las, é necessário declará-las primeiro. Se uma fila já existir com as mesmas propriedades, nada é feito; caso contrário, ela será criada ou atualizada.

Propriedades importantes das filas incluem: nome, persistência, exclusividade (somente para uma conexão) e exclusão automática quando não há consumidores ativos.

III. Ligações (Bindings)

Para que mensagens sejam encaminhadas corretamente, o exchange precisa estar ligado a uma ou mais filas. As ligações definem estas regras de roteamento. Caso não haja ligações, pode-se configurar o comportamento do broker ao receber mensagens — descartá-las ou encaminhá-las para uma fila de mensagens privadas.

IV. Consumidores (Consumers)

Os consumidores podem obter mensagens de duas formas: assinando (subscribe) ou buscando (pull). O método pull é mais pesado e não recomendado. Cada consumidor tem um identificador chamado consumerTag, usado para cancelar a assinatura.

O mecanismo de confirmação de mensagens funciona em dois modos:

  1. Confirmação automática: o broker remove a mensagem assim que ela é enviada.
  2. Confirmação explícita: o consumidor avisa o broker após processar a mensagem, garantindo entrega confiável. Se ocorrer falha antes da confirmação, a mensagem será redirecionada para outro consumidor.

V. Mensagens (Messages)

As mensagens possuem atributos como tipo de conteúdo (content-type) e payload (dados). Elas podem ser persistentes, embora isso impacte o desempenho.

VI. Conexões (Connections)

No AMQP, as conexões são persistentes e operam sobre TCP, com suporte a TLS.

VII. Canais (Channels)

Os canais são conexões leves compartilhadas dentro de uma conexão TCP. Cada canal possui um ID único e é isolado. Recomenda-se criar um canal por thread em sistemas multithreaded.

VIII. Hosts Virtuais (Virtual Hosts)

Um mesmo broker pode conter múltiplos hosts virtuais, permitindo ambientes isolados para diferentes aplicações.

IX. Código de exemplo Com base nos conceitos acima, o tipo de exchange mais utilizado é o Topic, pois oferece flexibilidade máxima. Ele permite usar padrões de roteamento com "." e caracteres curinga (* e #).

  1. Instalação via Docker: docker run -it --rm --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3.13-management Após iniciar, execute: rabbitmq-plugins enable rabbitmq_management Acesse http://localhost:15672 com usuário e senha guest/guest para gerenciar o RabbitMQ.
  2. Dependências:
<dependencies>
    <dependency>
        <groupId>com.rabbitmq</groupId>
        <artifactId>amqp-client</artifactId>
        <version>5.14.3</version>
    </dependency>
</dependencies>

  1. Código do produtor:
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class EmitLogTopic {

  private static final String EXCHANGE_NAME = "topic_logs";

  public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    try (Connection connection = factory.newConnection();
         Channel channel = connection.createChannel()) {

        channel.exchangeDeclare(EXCHANGE_NAME, "topic");

        String routingKey = getRouting(argv);
        String message = getMessage(argv);

        channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes("UTF-8"));
        System.out.println(" [x] Sent '" + routingKey + "':'" + message + "'");
    }
  }
}

Executar com parâmetros: "kern.critical" "A critical kernel error"

  1. Código do consumidor:
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;

public class ReceiveLogsTopic {

  private static final String EXCHANGE_NAME = "topic_logs";

  public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    Connection connection = factory.newConnection();
    Channel channel = connection.createChannel();

    channel.exchangeDeclare(EXCHANGE_NAME, "topic");
    String queueName = channel.queueDeclare().getQueue();

    if (argv.length < 1) {
        System.err.println("Usage: ReceiveLogsTopic [binding_key]...");
        System.exit(1);
    }

    for (String bindingKey : argv) {
        channel.queueBind(queueName, EXCHANGE_NAME, bindingKey);
    }

    System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

    DeliverCallback deliverCallback = (consumerTag, delivery) -> {
        String message = new String(delivery.getBody(), "UTF-8");
        System.out.println(" [x] Received '" +
            delivery.getEnvelope().getRoutingKey() + "':'" + message + "'");
    };
    channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { });
  }
}

Parâmetros válidos: "#" "kern.*" "*.critical" ou "kern.critical". Execute múltiplos consumidores para observar o comportamento.

Repositório: https://gitee.com/luxiao_gitee/rabbitmq.git

Tags: RabbitMQ AMQP mensagens fila exchange

Publicado em 8-23 14:30