How in RabbitMQ and PHP return task back to the queue?

Viewed 7267

How can I return message back to the queue if processing result did not suit me. Found only information about message acknowledgments but I think that it does not suit me. I need that if as a result of processing I get the parameter RETRY message is added back to the queue. And then this worker or another one picks it up again and tries to process it.

For example:

<?php
use PhpAmqpLib\Connection\AMQPStreamConnection;

echo ' [*] Waiting for messages. To exit press CTRL+C', "\n";

$connection = new AMQPStreamConnection($AMQP);
$channel = $connection->channel();

$channel->queue_declare('test', false, false, false, false);

$callback = function($msg) {
    $condition = json_decode($msg->body);

    if (!$condition) {
        # return to the queue
    }
};

$channel->basic_consume('test', '', false, true, false, false, $callback);

while(count($channel->callbacks)) {
    $channel->wait();
}

$channel->close();
$connection->close();
?>
3 Answers

This may be a bit late, but this is how you should do it with this version of php-amqplib "php-amqplib/php-amqplib": "^3.1" You need to set the no_ack parameter of basic consume method to false(which is the default) then explicitly specify it in the callback using the method nack on the AMQPMessage object passed to the callback

<?php
    use PhpAmqpLib\Connection\AMQPStreamConnection;
    
    echo ' [*] Waiting for messages. To exit press CTRL+C', "\n";
    
    $connection = new AMQPStreamConnection($AMQP);
    $channel = $connection->channel();
    
    $channel->queue_declare('test', false, false, false, false);
    
    $callback = function($msg) {
        $condition = json_decode($msg->body);
    
        if (!$condition) {
            // message will be added back to the queue
            $msg->nack(true);
        }
    };
    
    $channel->basic_consume('test', '', false, false, false, false, $callback);
    
    while(count($channel->callbacks)) {
        $channel->wait();
    }
    
    $channel->close();
    $connection->close();
    ?>
Related