Перейти к содержимому

Очередь сообщений - фоновая обработка задач в ядре

Выносим долгую операцию из запроса посетителя в фоновую обработку штатной очередью ядра. Разбираем класс сообщения, класс обработчика, брокер и режимы разбора самой очереди.

Механика

Штатная очередь ядра состоит из четырёх отдельных и независимых частей. Сообщение переносит данные, обработчик их принимает, брокер хранит очередь, а сама очередь связывает первое с вторым.

Сообщение обязано быть сериализуемым в простой и понятный обменный формат данных. В нём держат идентификаторы и короткие значения, а не объекты платформы и не замыкания.

Обработчик очереди - это отдельный класс ровно с одним рабочим методом обработки. Он получает готовое сообщение и делает работу, ради которой очередь и заводилась.

Очередь описывают в файле настроек продукта и связывают с её будущим обработчиком. Без этой записи отправка сообщения просто не находит, кому его вообще передать.

Брокер по умолчанию хранит всю очередь сообщений прямо в базе данных сайта. Это единственный поддержанный тип хранилища, и он обязан присутствовать в глобальных настройках.

Режим разбора очереди задают отдельной настройкой в файле настроек продукта. Разбор на обычных запросах включён по умолчанию, а для боевой среды берут консольный обработчик под присмотром супервизора.

Отправка сообщения в очередь не выполняет саму полезную работу немедленно. Запрос посетителя завершается сразу, а обработка происходит позже и отдельно от него.

Сообщение с одним лишь идентификатором внутри - довольно частая ошибка проектирования. Обработчик получит такое сообщение позже, а данные к тому моменту успеют измениться.

Секцию настроек очередей закрывают отдельным признаком доступа только для чтения файла. Это защищает описание очередей от случайной перезаписи через интерфейс настроек продукта.

Сам механизм пока довольно молодой и объявлен без всяких гарантий обратной совместимости. Версию ядра фиксируют, а критичные сценарии дополнительно дублируют проверкой результата обработки.

Очередь имеет смысл там, где посетителю не нужен готовый результат прямо сейчас. Выгрузка заказа в учётную систему, пересчёт остатков после обмена, обращение к чужому сервису - всё это спокойно живёт в фоне и не держит страницу.

А вот там, где результат нужен посетителю на экране, очередь только вредит. Ответ придётся ждать опросом состояния, и простая задача превращается в связку из трёх механизмов вместо одного обычного вызова.

Шаги

  1. Убедиться, что установленная версия главного модуля вообще поддерживает штатные очереди сообщений.
  2. Описать брокер сообщений по умолчанию в глобальных настройках всего продукта.
  3. Завести класс сообщения с простыми полями и понятными правилами его собственной сериализации.
  4. Написать класс обработчика с единственным методом обработки одного сообщения.
  5. Описать саму очередь и связать её с нужным обработчиком в файле настроек.
  6. Выбрать подходящий режим разбора очереди и запустить его под присмотром внешнего супервизора.

Код

Заводим класс сообщения:

use Bitrix\Main\Messenger\Entity\AbstractMessage;
use Bitrix\Main\Messenger\Entity\MessageInterface;
class OrderExportMessage extends AbstractMessage
{
public function __construct(
public readonly int $orderId, // простые поля сериализуются легко
public readonly string $reason
) {}
public function jsonSerialize(): mixed
{
return ['orderId' => $this->orderId, 'reason' => $this->reason];
}
public static function createFromData(array $data): MessageInterface
{
return new static(...$data);
}
}

Сообщение описывает задачу, а не тащит с собой её результат. Объекты заказа и корзины в него не кладут: они не переживут сериализации и устареют к моменту обработки.

Пишем обработчик:

use Bitrix\Main\Messenger\Receiver\AbstractReceiver;
class OrderExportReceiver extends AbstractReceiver
{
/** @param OrderExportMessage $message */
protected function process(MessageInterface $message): void
{
$this->exportOrder($message->orderId, $message->reason);
}
}
// исключение внутри обработки означает неудачную попытку, а не потерю сообщения
// журнал с идентификатором сообщения обязателен: разбор идёт в другом процессе

Обработчик занимается одной задачей и ничего не знает о том, кто её поставил. Такой класс легко проверить отдельным тестом, передав ему готовое сообщение прямо руками.

Описываем брокер и очередь:

/bitrix/.settings.php
'messenger' => ['value' => [
'run_mode' => 'cli', // web - разбор на обычных запросах
'brokers' => ['default' => ['type' => 'db']],
'queues' => [
'order_export' => ['handler' => \Vendor\Module\OrderExportReceiver::class],
],
], 'readonly' => true],

Брокер по умолчанию обязателен, даже когда очередей несколько. Признак только для чтения защищает это описание от перезаписи через интерфейс настроек.

Отправляем сообщение в очередь:

$message = new OrderExportMessage($order->getId(), 'статус изменён');
$message->send('order_export');
// запрос посетителя на этом заканчивается, работа уходит в фон
// имя очереди должно совпадать с описанным в настройках

Отправка занимает доли миллисекунды: сообщение просто ложится в хранилище. Именно поэтому очередь и ставят между быстрым ответом посетителю и медленной работой с чужой системой.

Запускаем разбор очереди:

Окно терминала
php /home/bitrix/www/bitrix.php messenger:consume
# в боевой среде процесс держит супервизор и перезапускает при падении
# режим разбора на обычных запросах годится для стенда, но не для нагрузки

Консольный обработчик разбирает очередь непрерывно и не зависит от посещаемости. Разбор на обычных запросах проще, но на нагруженном сайте отнимает время у посетителей.

Проверяем состояние очереди:

$rows = \Bitrix\Main\Messenger\Internals\MessengerMessageTable::getList([
'select' => ['ID', 'QUEUE', 'STATUS', 'DATE_CREATE'],
'order' => ['ID' => 'DESC'], 'limit' => 10,
])->fetchAll();
print_r($rows); // накопившиеся сообщения означают, что разбор не идёт
// статус сообщения показывает, взяли его в работу или ещё нет

Растущее число необработанных сообщений - первый признак остановившегося разбора. На боевом сайте такую проверку выносят в наблюдение вместе с очередью почтовых событий.

Отдельный вопрос - что именно считать успешной обработкой одного сообщения. Обработчик, который молча проглатывает ошибку чужого сервиса, выглядит работающим, а задача при этом не выполнена ни разу.

Поэтому у обработчика очереди заводят журнал с идентификатором сообщения и итогом работы. По нему видно и число повторов, и застрявшие задачи, и момент, когда чужая система начала отвечать отказом.

Ограничения

Сам механизм объявлен без всяких гарантий обратной совместимости своего API. Он появился недавно, и его интерфейсы могут поменяться в очередном обновлении продукта.

Хранилище очереди поддержано пока только одно - это самая обычная база данных. Внешние брокеры сообщений штатно не поддерживаются, и их подключение остаётся ручной работой.

Очередь сама по себе не заменяет надёжную доставку сообщений. Повторные попытки, порядок обработки и защита от дублей продумываются на стороне обработчика.

Фоновая обработка заметно усложняет разбор возникающих сбоев. Ошибка проявляется не в запросе посетителя, а позже и в другом процессе, поэтому журнал обработчика обязателен.

Типичные проблемы

Сообщение отправлено, а обработчик не вызывается.

Очередь не описана в настройках или указана в отправке не тем именем. Отправка кладёт сообщение в хранилище, но получателя для него так и нет.

Обработка падает на сериализации сообщения.

В сообщение положили объект платформы или замыкание с состоянием. Сообщение должно состоять только из простых значений и восстанавливаться из обычного массива.

Обработчик получил устаревшие данные.

В сообщении лежал только идентификатор, а данные успели измениться до момента обработки. Нужное состояние передают вместе с сообщением или перечитывают в обработчике осознанно.

На боевом сайте очередь разбирается рывками.

Включён разбор на обычных запросах: он зависит от посещаемости и отнимает время у посетителей. Для боевой нагрузки берут консольный обработчик очереди.

Настройки очередей затёрлись после правки в админке.

Секция настроек очередей не закрыта отдельным признаком доступа только для чтения. Сохранение настроек через интерфейс административной части переписывает этот файл целиком.

Частые вопросы

Чем очередь отличается от агента?

Агент выполняется по расписанию и ничего не знает о конкретной задаче. Очередь принимает именно задачу с данными и обрабатывает её отдельным классом, как только доходит очередь.

Что класть в сообщение?

Простые значения: идентификаторы, коды, короткие строки. Объекты платформы не переживут сериализации, а большие массивы раздувают хранилище очереди.

Как обрабатывать ошибки в обработчике?

Записывать в журнал и решать, повторять ли задачу. Логика повторов и защита от дублей - ответственность прикладного кода, а не механизма очереди.

Можно ли подключить внешний брокер сообщений?

Штатно поддержано только хранение в базе. Свой брокер технически возможен, но это ручная реализация со своей поддержкой.

Стоит ли переносить в очередь отправку писем?

Почтовые события и так уходят через свою очередь. Механизм очередей нужен для своей долгой работы: выгрузок, обращений к чужим системам, тяжёлых пересчётов.

Смежное

Первоисточник