RabbitMQ 等待多个队列完成

2022-01-11 00:00:00 queue rabbitmq messaging php symfony

好的,这里是正在发生的事情的概述:

Ok here is an overview of what's going on:

    M <-- Message with unique id of 1234
    |
    +-Start Queue
    |
    |
    | <-- Exchange
   /|
  / | 
 /  |   <-- bind to multiple queues
Q1  Q2  Q3
   |   / <-- start of the problem is here
   |  / 
   | /
   |/
    |
    Q4 <-- Queues 1,2 and 3 must finish first before Queue 4 can start
    |
    C <-- Consumer 

所以我有一个exchange,推送到多个队列,每个队列都有一个任务,一旦所有任务都完成了,才可以启动队列4.

So I have an exchange that pushes to multiple queues, each queue has a task, once all tasks are completed, only then can Queue 4 start.

因此,唯一 id 为 1234 的消息被发送到交换器,交换器将其路由到所有任务队列(Q1、Q2、Q3 等...),当消息 id 为 1234 的所有任务都完成后,运行消息 ID 1234 的 Q4.

So message with unique id of 1234 gets sent to the exchange, the exchange routes it to all the task queues ( Q1, Q2, Q3, etc... ), when all the tasks for message id 1234 have completed, run Q4 for message id 1234.

我该如何实现?

使用 Symfony2、RabbitMQBundle 和 RabbitMQ 3.x

Using Symfony2, RabbitMQBundle and RabbitMQ 3.x

资源:

  • http://www.rabbitmq.com/tutorials/amqp-concepts.html李>
  • http://www.rabbitmq.com/tutorials/tutorial-six-python.html

更新 #1

好的,我想这就是我要找的:

Ok I think this is what I'm looking for:

  • https://github.com/videlalvaro/Thumper/tree/master/examples/parallel_processing

带有并行处理的 RPC,但是如何将 Correlation Id 设置为我的唯一 ID 以对消息进行分组并识别哪个队列?

RPC with Parallel Processing, but how do I set the Correlation Id to be my unique id to group the messages and also identify what queue?

推荐答案

在RabbitMQ 站点上的 RPC 教程,有一种方法可以传递一个相关 id",可以识别您的消息给队列中的用户.

In the RPC tutorial at RabbitMQ's site, there is a way to pass around a 'Correlation id' that can identify your messages to users in the queue.

我建议在您的消息中使用某种 id 到前 3 个队列中,然后有另一个过程将消息从 3 个队列中取出到某种类型的存储桶中.当这些存储桶收到我假设完成的 3 个任务时,将最终消息发送到第 4 个队列进行处理.

I'd recommend using some sort of id with your messages into the first 3 queues and then have another process to dequeue messages from the 3 into buckets of some sort. When those buckets receive what I'm assuming is the completion of there 3 tasks, send the final message off to the 4th queue for processing.

如果您要为一个用户向每个队列发送超过 1 个工作项,您可能需要进行一些预处理以找出特定用户放入队列中的项目数量,以便在 4 之前出队的进程知道要多少排队前期待.

If you are sending more than 1 work item to each queue for one user, you might have to do a little preprocessing to find out how many items a particular user placed into the queue so the process dequeuing before 4 knows how many to expect before queuing up.

我用 C# 编写 rabbitmq,很抱歉我的伪代码不是 php 样式

I do my rabbitmq in C#, so sorry my pseudo code isn't in php style

// Client
byte[] body = new byte[size];
body[0] = uniqueUserId;
body[1] = howManyWorkItems;
body[2] = command;

// Setup your body here

Queue(body)

<小时>

// Server
// Process queue 1, 2, 3
Dequeue(message)

switch(message.body[2])
{
    // process however you see fit
}

processedMessages[message.body[0]]++;

if(processedMessages[message.body[0]] == message.body[1])
{
    // Send to queue 4
    Queue(newMessage)
}

<小时>

对更新 #1 的回应

与其将客户端视为终端,不如将客户端视为服务器上的进程.因此,如果您在 这个,那么您需要做的就是让服务器处理用户唯一 ID 的生成并将消息发送到适当的队列:

Instead of thinking of your client as a terminal, it might be useful to think of the client as a process on a server. So if you setup an RPC client on a server like this one, then all you need to do is have the server handle the generation of a unique id of a user and send the messages to the appropriate queues:

    public function call($uniqueUserId, $workItem) {
        $this->response = null;
        $this->corr_id = uniqid();

        $msg = new AMQPMessage(
            serialize(array($uniqueUserId, $workItem)),
            array('correlation_id' => $this->corr_id,
            'reply_to' => $this->callback_queue)
        );

        $this->channel->basic_publish($msg, '', 'rpc_queue');
        while(!$this->response) {
            $this->channel->wait();
        }

        // We assume that in the response we will get our id back
        return deserialize($this->response);
    }


$rpc = new Rpc();

// Get unique user information and work items here

// Pass even more information in here, like what queue to use or you could even loop over this to send all the work items to the queues they need.
$response = rpc->call($uniqueUserId, $workItem);

$responseBuckets[array[0]]++;

// Just like above code that sees if a bucket is full or not

相关文章