Conceptos Fundamentales sobre Exchanges
En el ecosistema de RabbitMQ, los mensajes no se entregan directamente a las colas. En su lugar, el productor envía mensajes a un Exchange (intercambiador), que se encarga de enrutar los datos hacia una o varias colas basándose en reglas específicas denominadas bindings.
Existen cuatro tipos principales de intercambiadores:
- Fanout (Abanico): Es el modelo más simple. Funciona mediante la difusión masiva (broadcast), enviando cada mensaje recibido a todas las colas que estén vinculadas a él, sin filtrar. Es el más eficiente en términos de velocidad de procesamiento.
- Direct (Directo): Utiliza una
routing_keypara entregar mensajes a colas específicas. Si la clave de enrutamiento del mensaje coincide exactamente con la clave de vinculación de la cola, el mensaje se entrega. Es ideal para sisteams con prioridades o categorías fijas. - Topic (Temas): Permite un enrutamiento flexible basado en patrones. Las claves de enrutamiento deben ser palabras separadas por puntos (ej.
"sistema.auth.error"). Se utilizan comodines:*: Sustituye exactamente una palabra.#: Sustituye cero o más palabras.
- Headers (Cabeceras): Ignora la clave de enrutamiento y utiliza los atributos del encabezado del mensaje. Al vincular una cola, se define un argumento
x-matchque puede serall(debe coincidir todo) oany(basta con que coincida un atributo).
Modelo de Colas de Trabajo (Work Queues)
Este patrón se utiliza para distribuir tareas pesadas entre varios consumidores. El objetivo es evitar el bloqueo por tareas intensivas mediante el reparto de carga. RabbitMQ utiliza por defecto un mecanismo de Round-robin, donde cada mensaje es procesado por un único consumidor, distribuyéndolos de forma equitativa entre todos los trabajadores conectados.
Modelo de Publicación/Suscripción
A diferencia de las colas de trabajo, donde un mensaje llega a un solo destinatario, en este modelo el mensaje se entrega a todos los interesados. Es la base de los sistemas dirigidos por eventos.
Implementación de un Emisor de Notificaciones (Productor)
public class EmisorLogs {
private final static String NOMBRE_EXCHANGE = "logs_difusion";
public static void main(String[] args) throws Exception {
ConnectionFactory con Factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection conexion = factory.newConnection();
Channel canal = conexion.createChannel()) {
// Definimos un exchange de tipo FANOUT
canal.exchangeDeclare(NOMBRE_EXCHANGE, BuiltinExchangeType.FANOUT, false);
String contenido = "Evento de sistema: Nueva conexión detectada";
// Publicamos al exchange sin especificar una routing key
canal.basicPublish(NOMBRE_EXCHANGE, "", null, contenido.getBytes("UTF-8"));
System.out.println(" [Sent] Notificación enviada: " + contenido);
}
}
}
Implementación de un Consumidor de Consola
public class MonitorConsola {
private final static String NOMBRE_EXCHANGE = "logs_difusion";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection conexion = factory.newConnection();
Channel canal = conexion.createChannel();
canal.exchangeDeclare(NOMBRE_EXCHANGE, BuiltinExchangeType.FANOUT, false);
// Creamos una cola temporal con nombre aleatorio
String nombreCola = canal.queueDeclare().getQueue();
canal.queueBind(nombreCola, NOMBRE_EXCHANGE, "");
DeliverCallback deliverCallback = (tag, mensaje) -> {
String raw = new String(mensaje.getBody(), "UTF-8");
System.out.println(" [Monitor] Recibido: " + raw);
};
canal.basicConsume(nombreCola, true, deliverCallback, tag -> {});
}
}
Patrón de Enrutamiento (Routing)
El enrutamiento permite suscribirse únicamente a un subconjunto de los mensajes. Por ejemplo, en un sistema de logs, podríamos tener un consumidor que solo guarde errores críticos en disco, mientras otro muestra todos los logs en pantalla.
Para esto, se utiliza el exchange de tipo Direct. La lógica de vinculación cambia de la siguiente manera:
// Vinculación para persistencia de errores graves
String colaErrores = "cola_critica";
canal.queueDeclare(colaErrores, true, false, false, null);
canal.queueBind(colaErrores, "EXCHANGE_DIRECTO", "ERROR");
// Vinculación para monitoreo general
String colaInfo = "cola_general";
canal.queueDeclare(colaInfo, true, false, false, null);
canal.queueBind(colaInfo, "EXCHANGE_DIRECTO", "INFO");
canal.queueBind(colaInfo, "EXCHANGE_DIRECTO", "WARNING");
canal.queueBind(colaInfo, "EXCHANGE_DIRECTO", "ERROR");
Al publicar, el emisor debe especificar la clave correspondiente:
canal.basicPublish("EXCHANGE_DIRECTO", "ERROR", null, "Fallo crítico en DB".getBytes());
Invocación de Métodos Remotos (RPC)
Aunque RabbitMQ es asíncrono por naturaleza, es posible implementar RPC (Remote Procedure Call). El cliente envía una petición con una propeidad replyTo (cola de respuesta) y un correlationId (ID único de mensaje). El servidor procesa la solicitud y devuelve el resultado a la cola especificada, permitiendo que el cliente asocie la respuesta con la petición original.