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.