Multiplexage du message RabbitMQ
Par exemple, j'ai 4 sources qui publient des métriques. Je voudrais multiplexer / fusionner tous ces messages dans une file d'attente / échange
--------+----+----+----+----+ -------+---------+----+---------+---------+
Source1 | M1 | M2 | M3 | | => Result | M1 | M4 | M2 | M3 | M6 | M5 | M7 |
Source2 | M4 | | | M5 |
Source3 | | | M6 | |
Source4 | | | | M7 |
For each queue:
* Read one message
* Publish message to the Result queue
Existe-t-il un moyen "natif" de le faire dans RabbitMQ ou dois-je écrire mon propre consommateur / éditeur?
MODIFIER 1
Un exemple à clarifier, disons qu'après un certain temps, j'ai
Processing "window"
+-+
Source1 |X|XXXXXXXXXXXXX
Source2 |Y|YYYYYYY
Source3 |Z|ZZZZZZZZZZ
Source4 |W|WW
+-+
Et puis plus tard
Processing "window"
+-+
Source1 XXX|X|XXXXXXXXXX
Source2 YYY|Y|YYYY
Source3 ZZZ|Z|ZZZZZZZ
Source4 WWW| |
+-+
Et puis plus tard
Processing "window"
+-+
Source1 XXXXXXXXX|X|XXXX
Source2 YYYYYYYY | |
Source3 ZZZZZZZZZ|Z|Z
Source4 WWW | |
+-+
La commande consommatrice résultante sera: X Y Z W X Y Z W X Y Z W X Y Z X Y Z X Y Z X Y Z X Y Z X Z X Z X Z X X X
X, Y, Z, W puis X, Y, Z, W puis X, Y, Z, W puis X, Y, Z ... X, Z ...
De cette façon, même si une source "envoie du spam", tous les autres messages provenant d'autres sources ont une chance d'être consommés.
Pour des raisons techniques / financières, je ne dois consommer qu'un seul message à la fois.
Le consommateur est beaucoup plus lent que les producteurs mais les producteurs publient beaucoup mais occasionnellement.
Si chaque source est publiée sur un échange lié à la même file d'attente, le résultat peut être XXXXXXXXXXXXXX YYYYYYYY ZZZZZZZZZZZ WWWou XXXXX Y XXXXX YYY XXX YYYY ZZZZZZZZZZZ WWW(selon le taux de publication de chaque source)
Réponses
Je pense que ce que vous voulez peut être réalisé simplement en exécutant un seul script qui s'abonne à toutes les files d'attente.
La condition clé est d'utiliser un seul thread d'application pour gérer tous les messages, quelle que soit la file d'attente dont ils proviennent. Ce à quoi cela ressemble varie en fonction de la langue et de la bibliothèque cliente que vous utilisez - si vous utilisez PHP, vous devrez vraiment faire tout votre possible pour ne pas être mono-thread, mais peut-être qu'il existe des bibliothèques clientes qui supposent que chaque rappel est sur un thread de travail distinct, et vous aurez besoin d'une ressource partagée pour qu'ils bloquent.
En ce qui concerne l'aspect réel de RabbitMQ, vous devrez:
- enregistrez un abonnement pour que le serveur vous envoie des messages, avec
basic.consume; c'est généralement recommandé plutôt que d'interroger explicitement avec debasic.gettoute façon - utiliser un seul «canal» pour tous les
basic.consumeappels - utiliser des accusés de réception manuels pour que les messages restent dans la file d'attente jusqu'à ce que votre processus soit terminé
- définir une limite de prélecture par file d'attente de 1 avec
basic.qos
Si vous avez 4 files d'attente, A, B, C et D, qui contiennent des quantités variables de messages lorsque vous démarrez le consommateur:
- Lorsque vous vous abonnez pour la première fois, la limite de prélecture signifie qu'un message de chaque file d'attente sera envoyé au canal; appelez-les A1, B1, C1 et D1
- La bibliothèque cliente lèvera un événement asynchrone dans votre application pour chacun d'entre eux à son tour
- Votre thread de travail unique gérera le premier de ces événements et commencera à traiter le message A1
- Jusqu'à ce que vous confirmiez manuellement ce message, aucun autre message ne peut arriver
- Une fois que vous avez accusé réception du premier message (A1), un nouveau message peut être pré-extrait de cette file d'attente (A2)
- Pendant ce temps, votre thread de travail débloquera et gérera le prochain événement qui a déjà été déclenché, pour le message B1
- Ce n'est qu'une fois que vous avez traité les événements en attente pour B1, C1 et D1 que le thread de travail verra l'événement pour le message A2
- Tant que les files d'attente ont des messages en attente, ils seront traités de manière circulaire. Même si toutes les files d'attente sauf une deviennent vides, elles reviendront en rotation dès qu'un message arrive, car un seul message de la file d'attente occupée aura été pré-extrait, le reste attendra simplement sur le serveur RabbitMQ.