Las colas de mensajes muertos (Dead Letter Queues - DLQ) en RabbitMQ son un mecanismo esencial para gestionar mensajes que no pueden ser procesados o entregados a los consumidores. Permiten redirigir estos mensajes a una cola separada para su posterior análisis o reintento, evitando así la pérdida de información y la saturación de las colas principales.
Concepto General de las Colas de Mensajes Muertos
Imaginemos un escenario de reserva de asientos para un tren. Un usuario realiza una solicitud de reserva, pero no completa el pago a tiempo. Si el asiento permaneciera reservado indefinidamente, otros usuarios no podrían adquirirlo. En tales casos, después de un período determinado, la reserva debería cancelarse y el asiento liberarse. Aquí es donde entra en juego la cola de mensajes muertos. Permite manejar estas reservas no pagadas o canceladas para un procesamiento posterior.
La implementación de este patrón implica la combinación de un Exchange de mensajes muertos y una Cola de mensajes muertos.
Cuando un mensaje en una cola normal es considerado "muerto" (no entregable), se reenvía a un exchange de mensajes muertos. Este exchange, a su vez, lo dirige a la cola de mensajes muertos designada, donde un consumidor específico se encargará de él.
Un mensaje se considera "muerto" en las siguientes situaciones:
- El mensaje es rechazado por el consumidor (usando
basic.rejectobasic.nack) con el parámetrorequeueestablecido enfalse. - El mensaje expira debido a un tiempo de vida (TTL) agotado sin ser consumido.
- La cola alcanza su longitud máxima permitida.
Causas de Generación de Mensajes Muertos
1. Rechazo de Mensajes por el Consumidor
Para implementar esto, se configuran un exchange y una cola de mensajes muertos, y se establecen los enlaces correspondientes.
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.beans.factory.annotation.Qualifier;
@Configuration
public class RabbitMqConfiguration {
// Definición de la cola principal
@Bean("primaryQueue")
public Queue primaryQueue() {
return QueueBuilder.durable("primary_queue_name")
.deadLetterExchange("dl_exchange") // Nombre del exchange de mensajes muertos
.deadLetterRoutingKey("dl_routing_key") // Clave de enrutamiento para la cola de mensajes muertos
.build();
}
// Definición del exchange de mensajes muertos
@Bean("deadLetterExchange")
public DirectExchange deadLetterExchange() {
return new DirectExchange("dl_exchange");
}
// Definición de la cola de mensajes muertos
@Bean("deadLetterQueue")
public Queue deadLetterQueue() {
return QueueBuilder.nonDurable("dl_queue_name").build();
}
// Enlace entre el exchange de mensajes muertos y la cola de mensajes muertos
@Bean
public Binding deadLetterBinding(@Qualifier("deadLetterExchange") DirectExchange dlExchange,
@Qualifier("deadLetterQueue") Queue dlQueue) {
return BindingBuilder.bind(dlQueue)
.to(dlExchange)
.with("dl_routing_key");
}
}
A continuación, se configuran los listeners para la cola principal y la cola de mensajes muertos.
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import com.example.messaging.model.MyMessage; // Asumiendo una clase de modelo
@Component
public class MessageListeners {
// Listener para la cola principal
@RabbitListener(queues = "primary_queue_name", ackMode = "MANUAL") // ackMode manual para control
public void receivePrimaryMessage(MyMessage message, com.rabbitmq.client.Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception {
System.out.println("Mensaje recibido en cola principal: " + message);
// Rechazar el mensaje y no devolverlo a la cola (se convierte en mensaje muerto)
channel.basicReject(tag, false);
}
// Listener para la cola de mensajes muertos
@RabbitListener(queues = "dl_queue_name")
public void receiveDeadLetterMessage(MyMessage message) {
System.out.println("Mensaje recibido en cola de mensajes muertos: " + message);
// Aquí se implementaría la lógica para manejar el mensaje muerto
}
}
Al rechazar un mensaje en la cola principal con requeue=false, este se enruta a la cola de mensajes muertos y es capturado por su listener.
2. Expiración de Mensajes (TTL)
Se configura un valor de Tiempo de Vida (TTL) para la cola principal.
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMqConfiguration {
// Definición de la cola principal con TTL
@Bean("primaryQueue")
public Queue primaryQueue() {
return QueueBuilder.durable("primary_queue_name")
.deadLetterExchange("dl_exchange")
.deadLetterRoutingKey("dl_routing_key")
.messageTTL(5000) // Tiempo de vida del mensaje en milisegundos (5 segundos)
.build();
}
// ... (otras configuraciones de exchange y cola de mensajes muertos) ...
}
Se desactiva temporalmente el listener de la cola principal para este escenario.
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import com.example.messaging.model.MyMessage;
@Component
public class MessageListeners {
// Listener principal comentado para este ejemplo
// @RabbitListener(queues = "primary_queue_name")
// public void receivePrimaryMessage(MyMessage message) {
// System.out.println("Mensaje recibido en cola principal: " + message);
// }
// Listener para la cola de mensajes muertos
@RabbitListener(queues = "dl_queue_name")
public void receiveDeadLetterMessage(MyMessage message) {
System.out.println("Mensaje recibido en cola de mensajes muertos (por expiración): " + message);
}
}
Si ningún consumidor procesa el mensaje dentro del TTL, este se moverá a la cola de mensajes muertos.
3. Cola Llena (Longitud Máxima)
Se establece un límite en la cantidad de mensajes que una cola puede contener.
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMqConfiguration {
// Definición de la cola principal con longitud máxima
@Bean("primaryQueue")
public Queue primaryQueue() {
return QueueBuilder.durable("primary_queue_name")
.deadLetterExchange("dl_exchange")
.deadLetterRoutingKey("dl_routing_key")
.maxQueueLength(3) // Máximo de 3 mensajes en la cola
.build();
}
// ... (otras configuraciones) ...
}
Similar al caso anterior, se puede comentar el listener de la cola principal.
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import com.example.messaging.model.MyMessage;
@Component
public class MessageListeners {
// Listener principal comentado
// @RabbitListener(queues = "primary_queue_name")
// public void receivePrimaryMessage(MyMessage message) {
// System.out.println("Mensaje recibido en cola principal: " + message);
// }
// Listener para la cola de mensajes muertos
@RabbitListener(queues = "dl_queue_name")
public void receiveDeadLetterMessage(MyMessage message) {
System.out.println("Mensaje recibido en cola de mensajes muertos (por longitud): " + message);
}
}
Cuando se envían más mensajes que la longitud máxima permitida, los mensajes más antiguos son desalojados y se convierten en mensajes muertos.
Flujo de Trabajo
Para que estas configuraciones funcionen, es crucial asegurarse de que la cola principal (con TTL o longitud máxima) o el exchange al que está enlazada, no existan previamente en el momento de la declaración si se desea que las nuevas configuraciones de TTL o longitud máxima se apliquen. De lo cnotrario, RabbitMQ puede mantener las propiedades de la cola existente.
Se utiliza un RabbitTemplate para enviar mensajes.
import org.junit.jupiter.api.Test;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import com.example.messaging.model.MyMessage;
@SpringBootTest
public class ProducerTest {
@Autowired
private RabbitTemplate rabbitTemplate;
@Test
void sendMessageToPrimaryQueue() {
MyMessage message = new MyMessage();
message.setContent("Datos del usuario");
// Asumiendo que la clave de enrutamiento es 'primary_key' y el exchange directo es 'primary_exchange'
rabbitTemplate.convertAndSend("primary_exchange", "primary_key", message);
}
@Test
void sendMultipleMessages() {
for (int i = 0; i < 5; i++) {
MyMessage message = new MyMessage();
message.setContent("Mensaje número " + i);
rabbitTemplate.convertAndSend("primary_exchange", "primary_key", message);
}
}
}
Tras enviar los mensajes y observar el comportamiento de los listeners, se verificará que los mensajes que cumplen las condiciones (rechazados, expirados o que llenan la cola) son procesados por el listener de la cola de mensajes muertos.
Entorno y Referencias
- JDK: 17.0.6
- Maven: 3.6.3
- SpringBoot: 3.0.4
- spring-boot-starter-amqp: 3.0.4
- jackson-databind: 2.14.2
Referencia: Video de referencia (Bilibili)