Estoy usando la función round robin de RabbitMQ para enviar mensajes entre varios consumidores, pero solo uno de ellos recibe el mensaje real a la vez.
Mi problema es que mis mensajes representan tareas y me gustaría tener sesiones locales (estado) en mis consumidores. Sé de antemano qué mensajes pertenecen a qué sesión, pero no sé cuál es la mejor manera (¿o hay alguna manera?) de hacer que RabbitMQ se envíe a los consumidores utilizando un algoritmo que especifique.
No quiero escribir mi propio servicio de orquestación porque se convertirá en un cuello de botella y no quiero que mis productores sepan qué consumidor tomará sus mensajes porque perderé el desacoplamiento que obtengo con Rabbit.
¿Hay alguna manera de hacer que RabbitMQ envíe mis mensajes a los consumidores en función de un algoritmo/regla predefinidos en lugar de turnos rotativos?
Aclaración: Uso varios microservicios escritos en diferentes idiomas y cada servicio tiene su propio trabajo. Me comunico entre ellos usando mensajes protobuf. Le doy a cada nuevo mensaje un UUID . Si un consumidor recibe un mensaje, puede crear un mensaje de respuesta a partir de él (es posible que esta no sea la terminología correcta, ya que los productores y los consumidores están desacoplados y no se conocen entre sí) y este UUID se copia en el nuevo mensaje. Esto forma una canalización de transformación de datos y este "proceso" se identifica mediante el UUID (el processId ). Mi problema es que es posible que tenga varios consumidores de trabajadores y necesito que un trabajador se ciña a un UUID si lo ha visto antes. Tengo esta necesidad porque
Dado que RabbitMQ distribuye tareas entre los trabajadores mediante turnos rotativos, no puedo forzar que mis procesos se adhieran a un trabajador. Tengo varias advertencias:
Si hay una solución alternativa que no implique cambiar el algoritmo de todos contra todos y no rompa mis restricciones, ¡también está bien!
Si no desea optar por un servicio de orquestación, puede probar una topología como esa:
En aras de la simplicidad, asumo que su processId de proceso se usa como la clave de enrutamiento (en el mundo real, es posible que desee almacenarlo en el encabezado y usar el intercambio de header en su lugar).
Un mensaje entrante será aceptado por el Intercambio entrante (tipo: directo), que tiene un atributo alternative-exchange establecido para apuntar al No Session Exchange (fanout).
Esto es lo que dicen los documentos de RabbitMQ sobre los 'Intercambios alternativos':
A veces es deseable permitir que los clientes manejen los mensajes que un intercambio no pudo enrutar (es decir, porque no había colas enlazadas o enlaces no coincidentes).
Ejemplos típicos de esto son
- detectar cuándo los clientes publican accidental o maliciosamente mensajes que no se pueden enrutar
- "o si no" semántica de enrutamiento donde algunos mensajes se manejan especialmente y el resto por un controlador genérico
La función Alternate Exchange ("AE") de RabbitMQ aborda estos casos de uso.
(estamos particularmente interesados en el caso de uso or else aquí)
Cada consumidor creará su propia cola y la vinculará al Intercambio entrante , utilizando processId(s) de proceso para las sesiones que conoce hasta el momento, como la clave de enrutamiento del enlace.
De esta forma, solo recibirá mensajes para las sesiones que le interesen.
Además, todos los consumidores se unirán a la Cola sin sesión compartida.
Si ingresa un mensaje con un processId de proceso previamente desconocido, no habrá un enlace específico para él registrado con el Intercambio entrante, por lo que se redirigirá al Intercambio sin sesión => Cola sin sesión y se enviará a uno de los Consumidores en una manera usual (round-robin).
Luego, un consumidor registrará un nuevo enlace para él con el Intercambio entrante (es decir, comenzará una nueva "sesión"), de modo que luego recibirá todos los mensajes posteriores con este processId de proceso.
Una vez finalizada la "sesión", deberá eliminar el enlace correspondiente (es decir, cerrar la "sesión").