Patrones de Mensajes con RabbitMQ: Mecanismos Avanzados

Control de Velocidad del Consumidor

Cuando el servidor de RabbitMQ acumula miles de mensajes sin procesar, enviarlos todos simultáneamente a un solo consumidor puede sobrecargarlo. La funcionalidad QoS (Calidad de Servicio) de RabbitMQ mitiga este problema. Al deshabilitar el reconocimiento automático (autoAck=false), se puede establecer un límite (prefetchCount) de mensajes no confirmados que un consumidor puede recibir a la vez. Si este límite se alcanza, el consumidor se bloquea hasta que confirme algún mensaje.

Implementación:

Primero, definimos un consumidor que confirme los mensajes manualmente después de procesarlos.

package com.example.rabbitmq.control;

import com.rabbitmq.client.*;
import java.io.IOException;

public class LimitConsumer extends DefaultConsumer {
    private final Channel channel;

    public LimitConsumer(Channel channel) {
        super(channel);
        this.channel = channel;
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
        String message = new String(body);
        System.out.println("[LimitConsumer] Mensaje recibido: " + message);
        try {
            // Simula un procesamiento que lleva tiempo
            Thread.sleep(500);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        // Confirma el mensaje de manera individual
        channel.basicAck(envelope.getDeliveryTag(), false);
    }
}

El código del consumidor principal establece el intercambio, la cola y el parámetro QoS. Es crucial configurar autoAck en false para que el control de velocidad funcione.

package com.example.rabbitmq.control;

import com.rabbitmq.client.*;
import java.util.concurrent.TimeoutException;

public class ConsumerApp {
    public static void main(String[] args) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {

            String exchange = "limit_control_exchange";
            String queue = "limit_control_queue";
            String bindingKey = "control.#";

            channel.exchangeDeclare(exchange, BuiltinExchangeType.TOPIC, true);
            channel.queueDeclare(queue, true, false, false, null);
            channel.queueBind(queue, exchange, bindingKey);

            // Habilitar control de velocidad: máximo 1 mensaje no confirmado por consumidor
            channel.basicQos(1);

            System.out.println("[ConsumerApp] Esperando mensajes. Para salir presione CTRL+C");
            channel.basicConsume(queue, false, new LimitConsumer(channel));

            // Mantener el consumidor activo
            System.in.read();
        }
    }
}

El productor envía múltiples mensajes al intercambio configurado.

package com.example.rabbitmq.control;

import com.rabbitmq.client.*;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeoutException;

public class ProducerApp {
    public static void main(String[] args) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {

            String exchange = "limit_control_exchange";
            String routingKey = "control.log";

            for (int i = 0; i < 10; i++) {
                String message = "Evento #" + (i + 1);
                channel.basicPublish(exchange, routingKey, null, message.getBytes(StandardCharsets.UTF_8));
                System.out.println("[ProducerApp] Enviado: " + message);
            }
        }
    }
}

Confirmación Manual (ACK/NACK) y Reencolado de Mensajes

La confirmación automática puede llevar a la pérdida de mensajes si un consumidor falla durante el procesamiento. La confirmación manual ofrece un control preciso. Un consumidor puede confirmar (basicAck) un mensaje procesado exitosamente o rechazarlo (basicNack o basicReject). El parámetro requeue en basicNack determina si el mensaje debe ser devuelto a su cola original para ser reintentado.

Ejemplo de Lógica de Reintento:

Este consumidor intenta procesar un mensaje. Si el encabezado del mensaje indica que debe reintentarse (por ejemplo, si el conteo es 0), lo rechaza y lo reencola. De lo contrario, lo confirma.

package com.example.rabbitmq.ack;

import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.Map;

public class RetryConsumer extends DefaultConsumer {
    private final Channel channel;

    public RetryConsumer(Channel channel) {
        super(channel);
        this.channel = channel;
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
        Map<string object=""> headers = properties.getHeaders();
        Integer retryCount = (Integer) headers.get("retryCount");
        String message = new String(body);

        System.out.println("[RetryConsumer] Mensaje: " + message + ", Reintentos: " + retryCount);

        try {
            if (retryCount == null || retryCount < 3) {
                // Lanza una excepción simulada para forzar el reintento
                throw new RuntimeException("Error de procesamiento simulado");
            }
            // Procesamiento exitoso
            System.out.println("[RetryConsumer] Mensaje procesado exitosamente.");
            channel.basicAck(envelope.getDeliveryTag(), false);
        } catch (Exception e) {
            System.err.println("[RetryConsumer] Error: " + e.getMessage() + ". Reencolando mensaje.");
            // Rechaza el mensaje y lo reencola para reintento.
            channel.basicNack(envelope.getDeliveryTag(), false, true);
        }
    }
}</string>

El consumidor principal se suscribe a la cola con la lógica de reintento.

package com.example.rabbitmq.ack;

import com.rabbitmq.client.*;

public class ConsumerApp {
    public static void main(String[] args) throws Exception {
        // ... (Configuración de conexión similar al ejemplo anterior)
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {

            String exchange = "retry_exchange";
            String queue = "retry_queue";
            channel.exchangeDeclare(exchange, BuiltinExchangeType.DIRECT, true);
            channel.queueDeclare(queue, true, false, false, null);
            channel.queueBind(queue, exchange, "process");

            System.out.println("[ConsumerApp] Iniciado con lógica de reintentos.");
            channel.basicConsume(queue, false, new RetryConsumer(channel));
            System.in.read();
        }
    }
}

Mensajes con Tiempo de Vida (TTL)

RabbitMQ permite definir un tiempo de vida (TTL) para los mensajes, tanto a nivel de cola como individualmente en cada mensaje. Si un mensaje permanece en una cola por más tiempo que su TTL configurado, se convierte en un "dead letter" (carta muerta) y puede ser manejado por un intercambio de cartas muertas (DLX).

Ejemplo de Configuración por Cola:

Al declarar una cola, se puede establecer el argumento x-message-ttl en milisegundos.

Map<String, Object> args = new HashMap<>();
args.put("x-message-ttl", 60000); // Los mensajes expiran después de 60 segundos
channel.queueDeclare("ttl_queue", true, false, false, args);

Ejemplo de Configuración por Mensaje:

Al publicar un mensaje, se establece la propiedad expiration.

AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
        .expiration("10000") // TTL de 10 segundos
        .build();
channel.basicPublish("", "ttl_queue", props, "Mensaje temporal".getBytes());

Cola de Cartas Muertas (Dead Letter Queue - DLQ)

Una DLQ es una cola normal que actúa como destino para mensajes que no pudieron ser procesados (cartas muertas). Un mensaje se convierte en carta muerta si:

  • Es rechazado con requeue=false.
  • Su TTL expira.
  • La cola a la que pertenece alcanza su límite máximo de longitud (x-max-length).

Para configurar una DLQ, se asocia un intercambio de cartas muertas (DLX) a una cola a través de argumentos. Cuando un mensaje muere en esa cola, RabbitMQ lo publica automáticamente en el DLX.

Configuración Paso a Paso:

package com.example.rabbitmq.dlx;

import com.rabbitmq.client.*;
import java.util.HashMap;
import java.util.Map;

public class Setup {
    public static void main(String[] args) throws Exception {
        // ... (Configuración de conexión)
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {

            // 1. Declarar el intercambio y cola para cartas muertas
            String dlxExchange = "dead_letter_exchange";
            String dlxQueue = "dead_letter_queue";
            channel.exchangeDeclare(dlxExchange, BuiltinExchangeType.TOPIC, true);
            channel.queueDeclare(dlxQueue, true, false, false, null);
            channel.queueBind(dlxQueue, dlxExchange, "#");

            // 2. Declarar la cola principal con el argumento DLX
            String mainExchange = "main_processing_exchange";
            String mainQueue = "main_processing_queue";
            Map<String, Object> mainArgs = new HashMap<>();
            mainArgs.put("x-dead-letter-exchange", dlxExchange);
            mainArgs.put("x-message-ttl", 30000); // Opcional: TTL de 30s

            channel.exchangeDeclare(mainExchange, BuiltinExchangeType.DIRECT, true);
            channel.queueDeclare(mainQueue, true, false, false, mainArgs);
            channel.queueBind(mainQueue, mainExchange, "process");

            System.out.println("Infraestructura de DLQ configurada.");
            System.out.println("Los mensajes fallidos de 'main_processing_queue' irán a 'dead_letter_queue'.");
        }
    }
}

Un consumidor normal se conecta a main_processing_queue. Si falla, el mensaje eventualmente llegará a dead_letter_queue, donde otro consumidor (por ejemplo, uno de monitoreo o reintento) puede procesarlo.

Etiquetas: RabbitMQ patrones-de-mensajes qos TTL cola-de-cartas-muertas

Publicado el 7-24 21:13