Pesan Multiplexing RabbitMQ
Misalnya saya punya 4 sumber yang menerbitkan meterics. Saya ingin menggandakan / menggabungkan semua pesan ini dalam satu antrian / pertukaran
--------+----+----+----+----+ -------+---------+----+---------+---------+
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
Apakah ada cara "asli" untuk melakukan ini di RabbitMQ atau haruskah saya menulis Konsumen / Penerbit saya sendiri?
EDIT 1
Beberapa contoh untuk memperjelas, katakanlah setelah beberapa waktu saya punya
Processing "window"
+-+
Source1 |X|XXXXXXXXXXXXX
Source2 |Y|YYYYYYY
Source3 |Z|ZZZZZZZZZZ
Source4 |W|WW
+-+
Dan kemudian
Processing "window"
+-+
Source1 XXX|X|XXXXXXXXXX
Source2 YYY|Y|YYYY
Source3 ZZZ|Z|ZZZZZZZ
Source4 WWW| |
+-+
Dan kemudian
Processing "window"
+-+
Source1 XXXXXXXXX|X|XXXX
Source2 YYYYYYYY | |
Source3 ZZZZZZZZZ|Z|Z
Source4 WWW | |
+-+
Urutan konsumsi hasil adalah: 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 lalu X, Y, Z, W lalu X, Y, Z, W lalu X, Y, Z ... X, Z ...
Dengan cara ini, meskipun suatu sumber "melakukan spamming", semua pesan lain dari sumber lain memiliki kesempatan untuk dikonsumsi.
Untuk alasan teknis / keuangan, saya hanya perlu mengonsumsi 1 pesan setiap kali.
Konsumen jauh lebih lambat daripada produsen tetapi produsen menerbitkan banyak tetapi sesekali.
Jika setiap sumber dipublikasikan ke bursa yang terikat pada antrean yang sama, hasilnya mungkin XXXXXXXXXXXXXX YYYYYYYY ZZZZZZZZZZZ WWWatau XXXXX Y XXXXX YYY XXX YYYY ZZZZZZZZZZZ WWW(bergantung pada tingkat publikasi setiap sumber)
Jawaban
Saya pikir apa yang Anda inginkan dapat dicapai hanya dengan menjalankan satu skrip yang berlangganan semua antrian.
Persyaratan utamanya adalah menggunakan utas aplikasi tunggal untuk menangani semua pesan, terlepas dari antrian mana pesan tersebut datang. Seperti apa tampilannya akan bervariasi tergantung pada bahasa apa dan pustaka klien yang Anda gunakan - jika Anda menggunakan PHP, Anda harus benar-benar berusaha agar tidak menjadi utas tunggal, tetapi mungkin ada beberapa pustaka klien yang mengasumsikan setiap callback ada di thread pekerja terpisah, dan Anda memerlukan beberapa sumber daya bersama untuk memblokirnya.
Dalam hal sisi RabbitMQ yang sebenarnya, Anda perlu:
- mendaftarkan langganan ke server untuk mengirim pesan kepada Anda, dengan
basic.consume; ini umumnya direkomendasikan daripada polling secara eksplisit denganbasic.getanyway - gunakan satu "saluran" untuk semua
basic.consumepanggilan - gunakan pengakuan manual sehingga pesan tetap dalam antrian sampai proses Anda selesai
- setel batas prefetch per antrian 1 dengan
basic.qos
Jika Anda memiliki 4 antrian, A, B, C, dan D, yang memiliki jumlah pesan yang berbeda-beda saat Anda memulai konsumen:
- Saat Anda pertama kali berlangganan, batas prefetch berarti satu pesan dari setiap antrian akan dikirim ke saluran; sebut saja A1, B1, C1, dan D1
- Pustaka klien akan memunculkan peristiwa asinkron di aplikasi Anda untuk masing-masing peristiwa ini secara bergantian
- Rangkaian pekerja tunggal Anda akan menangani kejadian pertama ini, dan mulai memproses pesan A1
- Sampai Anda mengetahui pesan itu secara manual, tidak ada pesan lain yang bisa sampai
- Setelah Anda menerima pesan pertama (A1), pesan baru dapat diambil sebelumnya dari antrian itu (A2)
- Sementara itu, utas pekerja Anda akan membuka blokir dan menangani acara berikutnya yang sudah dimunculkan, untuk pesan B1
- Hanya setelah Anda memproses peristiwa tertunda untuk B1, C1, dan D1, thread pekerja akan melihat peristiwa untuk pesan A2
- Selama antrian memiliki pesan yang menunggu, maka akan diproses dengan cara round-robin. Bahkan jika semua kecuali satu dari antrian menjadi kosong, mereka akan masuk kembali ke rotasi segera setelah pesan tiba, karena hanya satu pesan dari antrian sibuk yang akan diambil sebelumnya, sisanya hanya akan menunggu di server RabbitMQ.