在 PHP 中使用 RabbitMQ
使用 RabbitMQ 发布和消费消息。
在 PHP 中使用 RabbitMQ 是 CoddyKit 上的免费 PHP Academy 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 PHP Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 PHP Academy 课程共包含 4 节课。
PHP 中的 RabbitMQ
RabbitMQ 是一个使用 AMQP 0-9-1 协议的代理。在 PHP 中,规范客户端是 php-amqplib/php-amqplib(纯 PHP),或者使用 C 支持的 ext-amqp。本课使用 php-amqplib,因为只要 Composer 能运行,它就能在任何环境中运行。
AMQP 核心模型包含三类参与者:交换器接收消息,绑定根据键路由消息,队列为消费者保存消息。掌握这一点后,其余内容只是细节。
安装客户端
使用 Composer 添加该库。它需要 sockets 和 bcmath 扩展,这两个扩展在 CLI PHP 构建中都很常见。
composer require php-amqplib/php-amqplib
# Connection target, e.g. amqp://guest:guest@localhost:5672/交换器/队列/绑定模型
生产者将消息发布到交换器,而不是直接发布到队列。交换器类型决定路由方式:
direct——路由键必须完全匹配。topic——支持类似order.*.eu的通配符模式。fanout——广播到所有已绑定的队列。headers——根据标头属性进行匹配。
绑定使用路由模式将队列连接到交换器。将生产者与队列拓扑解耦,正是交换器层存在的全部意义。
连接与声明
打开连接,获取通道,并声明拓扑。durable: true 可使 the 交换机/队列在消息代理重启后继续存在。声明具有幂等性——使用匹配参数声明已存在的实体时不会执行任何操作。
<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->exchange_declare('orders', 'topic', false, true, false);
$ch->queue_declare('orders.email', false, true, false, false);
$ch->queue_bind('orders.email', 'orders', 'order.created');
echo "Topology ready\n";
$ch->close();
$conn->close();发布消息
将 the 消息体封装在 AMQPMessage 中。设置 delivery_mode = 2 使 the 消息持久化——持久队列加持久化消息才能在重启后保留(只有其中一项并不足够)。
<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();
$payload = json_encode(['orderId' => 42, 'total' => 19.90]);
$msg = new AMQPMessage($payload, [
'content_type' => 'application/json',
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
'message_id' => bin2hex(random_bytes(8)),
]);
$ch->basic_publish($msg, 'orders', 'order.created');
echo "Published\n";
$ch->close();
$conn->close();使用手动确认消费
使用 basic_consume 注册回调。传入 no_ack = false,以便由您控制确认。仅在工作成功后确认;失败时,使用允许重新入队的 basic_nack,这样 RabbitMQ 会重新投递消息。
<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->basic_qos(null, 10, null); // prefetch 10
$callback = function ($msg) {
$data = json_decode($msg->getBody(), true);
try {
// ... do work ...
$msg->ack();
} catch (\Throwable $e) {
$msg->nack(true); // requeue
}
};
$ch->basic_consume('orders.email', '', false, false, false, false, $callback);
while ($ch->is_consuming()) {
$ch->wait();
}预取与公平分发
默认情况下,RabbitMQ 会以轮询方式向消费者分发消息,而不考虑每个消费者有多忙——较慢的消费者会不断积压工作。basic_qos(null, prefetch, null) 会限制消费者可以持有的未确认消息数量。
对于繁重且负载不均的任务,请将预取数量设置为较小的数值(例如 1–10),这样消息代理只会向有空闲容量的消费者发送新工作。这就是公平分发。
发布者确认
basic_publish 会立即返回,并不会告诉您 the 消息代理是否已接受 the 消息。要实现可靠发布,请启用发布者确认:the 消息代理会在消息安全持久化或路由完成后发送确认。
<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->confirm_select(); // enable confirms on this channel
$ch->set_ack_handler(fn($m) => print("confirmed\n"));
$ch->set_nack_handler(fn($m) => print("REJECTED\n"));
$ch->basic_publish(new AMQPMessage('hi'), 'orders', 'order.created');
$ch->wait_for_pending_acks(5.0); // block until confirmed or timeout死信处理
配置队列的 x-dead-letter-exchange 参数,使被拒绝的消息(不重新入队的 nack)或已过期消息路由到 DLX。对于仲裁队列,可结合 x-delivery-limit 自动限制重试次数。
<?php
use PhpAmqpLib\Wire\AMQPTable;
$args = new AMQPTable([
'x-dead-letter-exchange' => 'orders.dlx',
'x-dead-letter-routing-key' => 'order.failed',
'x-message-ttl' => 60000, // ms before expiry
]);
// false=passive, true=durable, false=exclusive, false=autodelete, args
$ch->queue_declare('orders.email', false, true, false, false, false, $args);
$ch->queue_declare('orders.dead', false, true, false, false);
$ch->queue_bind('orders.dead', 'orders.dlx', 'order.failed');在生产环境运行工作进程
针对 PHP RabbitMQ 工作进程的一些实践经验规则:
- 将消费者作为长期运行的 CLI 进程运行,并置于进程监管器(systemd / Supervisor)管理之下,使其在退出时自动重启。
- PHP 会随时间泄漏内存——在处理 N 条消息后或达到内存阈值时重启工作进程。
- 发送 AMQP 心跳,并处理
SIGTERM以实现平滑关闭(完成正在处理的消息,然后停止)。 - 使用仲裁队列实现 HA,而不是使用旧式镜像队列。
平滑关闭
部署会发送 SIGTERM。粗糙的工作进程会在消息处理到一半时退出,导致消息必须重新投递。安装信号处理程序,使其切换一个标志;完成当前消息并确认,然后干净地跳出消费循环。pcntl_async_signals(true) 允许 PHP 在等待 AMQP 的间隙传递信号。
<?php
pcntl_async_signals(true);
$running = true;
pcntl_signal(SIGTERM, function () use (&$running) {
$running = false; // stop after the current message
echo "SIGTERM: draining...\n";
});
while ($running && $ch->is_consuming()) {
try {
$ch->wait(null, false, 5); // wakes for signals
} catch (\PhpAmqpLib\Exception\AMQPTimeoutException $e) {
// idle tick - loop and re-check $running
}
}
$ch->close();
echo "Stopped cleanly\n";快速检查
在消息代理重启后仍然保留。
回顾
现在您可以用 PHP 构建真正的 RabbitMQ 工作流:
- 发布到交换机;通过绑定路由到队列。
- 持久队列 + 持久化消息可在重启后保留。
- 使用手动确认消费,并调节basic_qos 预取以实现公平分发。
- 使用发布者确认确保发送成功,并使用DLX处理毒消息。
- 在进程监管器下运行工作进程,并启用心跳和实现平滑关闭。
接下来学习 Kafka,适用于需要高吞吐流处理而非任务排队的场景。
常见问题解答
「在 PHP 中使用 RabbitMQ」课时是免费的吗?
是的 — 「在 PHP 中使用 RabbitMQ」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 PHP Academy 课程的其余内容,请升级到 CoddyKit PRO。 PHP Academy 课程共包含 4 节课。
「在 PHP 中使用 RabbitMQ」这节课中我会学到什么?
使用 RabbitMQ 发布和消费消息。 你通过在浏览器中直接运行的动手代码来练习 PHP Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 PHP Academy 需要有经验吗?
无需任何先前经验。CoddyKit 上的 PHP Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。
「在 PHP 中使用 RabbitMQ」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 PHP Academy 课中编写并运行代码吗?
能。每节 PHP Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- 异步消息为何重要
- 在 PHP 中使用 RabbitMQ
- 使用 PHP 的 Apache Kafka
- 构建事件驱动的工作流