Работа с RabbitMQ в PHP
Публикуйте и потребляйте сообщения с помощью RabbitMQ
«Работа с RabbitMQ в PHP» — бесплатный урок PHP Academy на CoddyKit. Это урок 2 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения PHP Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс PHP Academy содержит 4 уроков всего.
RabbitMQ в PHP
RabbitMQ — это брокер, поддерживающий протокол AMQP 0-9-1. В PHP стандартным клиентом служит php-amqplib/php-amqplib (чистый PHP) или расширение на основе C ext-amqp. В этом уроке используется php-amqplib, поскольку оно работает везде, где работает Composer.
Основная модель AMQP включает три участника: обменники получают сообщения, связи направляют их по ключу, а очереди хранят их для потребителей. Освойте эту модель — остальное лишь детали.
Установка клиента
Добавьте библиотеку с помощью Composer. Ей требуются расширения sockets и bcmath, оба обычно доступны в сборках PHP для CLI.
composer require php-amqplib/php-amqplib
# Connection target, e.g. amqp://guest:guest@localhost:5672/Модель обменника, очереди и связи
Производители публикуют сообщения в обменник, а не напрямую в очередь. Тип обменника определяет маршрутизацию:
direct— точное совпадение ключа маршрутизации.topic— шаблоны с подстановочными знаками, напримерorder.*.eu.fanout— рассылка во все связанные очереди.headers— сопоставление по атрибутам заголовков.
Связь соединяет очередь с обменником с помощью шаблона маршрутизации. Разделение производителей и топологии очередей — главная цель уровня обменников.
Подключение и объявление
Откройте соединение, получите канал и объявите топологию. durable: true позволяет обменнику и очереди пережить перезапуск брокера. Объявления идемпотентны: повторное объявление уже существующей сущности с совпадающими аргументами не приводит к действиям.
<?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();Публикация сообщения
Оберните тело в AMQPMessage. Установите delivery_mode = 2, чтобы сделать сообщение сохраняемым: именно устойчивая очередь и сохраняемое сообщение вместе позволяют пережить перезапуск, а по отдельности они этого не обеспечивают.
<?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 сразу возвращает управление и не сообщает, принял ли брокер сообщение. Для гарантированной публикации включите подтверждения публикации: брокер отправит подтверждение, когда сообщение будет надёжно сохранено или маршрутизировано.
<?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";Быстрая проверка
Устойчивость к перезапуску брокера.
Повторение
Теперь Вы умеете строить настоящий конвейер RabbitMQ на PHP:
- Публикуйте сообщения в обменник и маршрутизируйте их через привязки в очереди.
- Устойчивая очередь и сохраняемое сообщение вместе позволяют пережить перезапуск.
- Получайте сообщения с помощью ручного подтверждения и настраивайте предварительную выборку basic_qos для справедливого распределения.
- Используйте подтверждения публикации для гарантированной отправки и DLX для проблемных сообщений.
- Запускайте обработчики под управлением супервизора, с пульсами и корректным завершением.
Далее: Kafka, когда вместо постановки задач в очередь Вам нужна потоковая обработка с высокой пропускной способностью.
Часто задаваемые вопросы
Урок «Работа с RabbitMQ в PHP» бесплатный?
Да — полный текст урока «Работа с RabbitMQ в PHP» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс PHP Academy, подпишись на CoddyKit PRO. Курс PHP Academy содержит 4 уроков всего.
Чему я научусь в уроке «Работа с RabbitMQ в PHP»?
Публикуйте и потребляйте сообщения с помощью RabbitMQ Ты практикуешь PHP Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать PHP Academy?
Предыдущий опыт не требуется. PHP Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 2 из 4.
Сколько времени занимает урок «Работа с RabbitMQ в PHP»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке PHP Academy?
Да. Каждый урок PHP Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Зачем нужны асинхронные сообщения
- Работа с RabbitMQ в PHP
- Apache Kafka с PHP
- Создание событийно-ориентированных рабочих процессов