videlalvaro/php-amqplib
Pure PHP AMQP 0-9-1 client for RabbitMQ. Provides connections and channels for publishing and consuming messages, plus RabbitMQ extensions like publisher confirms, basic nack, consumer cancel notifications, and exchange-to-exchange bindings.
Installation
composer require videlalvaro/php-amqplib
Ensure ext-amqp or ext-php-amqplib is installed (if using PHP 8+ with php-amqplib extension).
Basic Connection
use PhpAmqpLib\Connection\AMQPStreamConnection;
$connection = new AMQPStreamConnection(
'localhost', // Host
5672, // Port
'guest', // Username
'guest' // Password
);
First Use Case: Publish a Message
$channel = $connection->channel();
$channel->queue_declare('hello', false, true, false, false);
$channel->basic_publish(
new \PhpAmqpLib\Message\AMQPMessage('Hello World!'),
'',
'hello'
);
Consume a Message
$channel->queue_declare('hello', false, true, false, false);
$channel->basic_consume('hello', '', false, true, false, false, function ($msg) {
echo "Received: " . $msg->body . "\n";
$msg->ack();
});
while ($channel->is_consuming()) {
$channel->wait();
}
examples/ folder in the repo for real-world use cases (e.g., RPC, pub/sub).PhpAmqpLib\Message\AMQPMessage for message construction.PhpAmqpLib\Channel\AMQPChannel for queue/exchange operations.// Producer (Publisher)
$channel->exchange_declare('logs', 'fanout', false, false, false);
$channel->basic_publish(
new AMQPMessage('Log entry: ' . time()),
'logs'
);
// Consumer (Subscriber)
$channel->queue_declare('logs_queue', false, true, false, false);
$channel->queue_bind('logs_queue', 'logs');
$channel->basic_consume('logs_queue', '', false, true, false, false, $callback);
// Worker (Consumer)
$channel->queue_declare('task_queue', false, true, false, false);
$channel->basic_qos(null, 1, null); // Fair dispatch (1 message at a time)
$channel->basic_consume('task_queue', '', false, true, false, false, $callback);
// Client
$channel->queue_declare('rpc_queue', false, true, false, false);
$corr_id = uniqid();
$channel->basic_publish(
new AMQPMessage(json_encode(['action' => 'add', 'params' => [2, 3]]), $corr_id),
'',
'rpc_queue'
);
// Server
$channel->queue_declare('rpc_queue', false, true, false, false);
$channel->basic_consume('rpc_queue', '', false, true, false, false, function ($msg) {
$data = json_decode($msg->body, true);
$response = ['result' => $data['params'][0] + $data['params'][1]];
$msg->reply_to = $msg->reply_to ?: 'rpc_queue';
$msg->correlation_id = $msg->correlation_id;
$channel->basic_publish(
new AMQPMessage(json_encode($response), $msg->correlation_id),
'',
$msg->reply_to
);
$msg->ack();
});
$channel->queue_declare('main_queue', false, true, false, false);
$channel->queue_declare('dlq', false, true, false, false);
$channel->queue_bind('dlq', '', 'main_queue', false, ['x-dead-letter-exchange' => '']);
$channel->basic_publish(
new AMQPMessage('Failing message', ['delivery_mode' => 2]),
'',
'main_queue'
);
Laravel Integration:
Use php-amqplib with Laravel’s Service Providers and Queues:
// config/queues.php
'connections' => [
'rabbitmq' => [
'driver' => 'rabbitmq',
'host' => env('RABBITMQ_HOST', 'localhost'),
'port' => env('RABBITMQ_PORT', 5672),
'user' => env('RABBITMQ_USER', 'guest'),
'password' => env('RABBITMQ_PASSWORD', 'guest'),
'vhost' => env('RABBITMQ_VHOST', '/'),
],
];
Extend Illuminate\Queue\QueueManager to use php-amqplib:
$this->app->extend('queue', function ($app) {
return new RabbitMQQueue(
new AMQPStreamConnection(...),
$app['queue.worker']
);
});
Retry Logic: Implement exponential backoff for failed messages:
$attempts = 0;
$max_attempts = 3;
$delay = 1000; // ms
$channel->basic_consume('queue', '', false, true, false, false, function ($msg) use (&$attempts, $max_attempts, $delay, $channel) {
try {
// Process message
$msg->ack();
$attempts = 0;
} catch (\Exception $e) {
if ($attempts < $max_attempts) {
$attempts++;
$channel->basic_nack($msg->delivery_info['delivery_tag'], false, true);
sleep($delay * $attempts);
} else {
$msg->nack(false); // Move to DLQ
}
}
});
Connection Management: Use connection pooling or reconnect logic:
$connection = new AMQPStreamConnection($host, $port, $user, $pass);
$retry = 3;
while ($retry--) {
try {
$channel = $connection->channel();
break;
} catch (\Exception $e) {
if ($retry === 0) throw $e;
sleep(1);
}
}
Connection Leaks:
__destruct() or a connection manager:
$channel->close();
$connection->close();
Message Acknowledgment (ACK/NACK) Misuse:
$msg->ack()) causes unacked messages to pile up.ack() or nack() messages in consumers.Poison Pills:
Race Conditions in RPC:
uniqid() or UUIDs for correlation_id.Durable vs. Non-Durable Queues:
false, true for durable queues:
$channel->queue_declare('queue', false, true, false, false);
Memory Limits:
$connection = new AMQPStreamConnection(
$host, $port, $user, $pass,
'/', // vhost
false, // login_method
false, // login_response
null, // locale
'AMQPLAIN', // connection_type
null, // user_context
0, // read_timeout
0, // write_timeout
null, // heartbeat
false, // keepalive
null, // channel_max
null, // frame_max
null, // heartbeat
null // channel_rpc_timeout
);
$connection->setReadWriteTimeout(15); // Adjust timeouts
How can I help you explore Laravel packages today?