Messenger — Symfony
Частину підрозділів ще не перекладено — вони показані англійською нижче в тексті або лишились в оригіналі. Готові фрагменти вже перевірені редактором.
Messenger надає шину повідомлень (message bus) із можливістю надсилати повідомлення й потім обробляти їх одразу у вашому застосунку або передавати їх через транспорти (наприклад, черги), щоб обробити пізніше. Компонент значною мірою натхненний серією блогових дописів про командні шини (blog posts about command buses) Матіаса Нобака та проєктом SimpleBus.
Встановлення
У застосунках, що використовують Symfony Flex, виконайте цю команду, щоб встановити messenger:
composer require symfony/messenger
Концепції
Перед використанням Messenger корисно ознайомитися з його основними концепціями та тим, як вони повʼязані між собою:
Sender (відправник): : Відповідає за серіалізацію та надсилання повідомлень кудись. Цим «кудись» може бути, наприклад, брокер повідомлень або стороннє API.
Receiver (отримувач): : Відповідає за отримання, десеріалізацію та передавання повідомлень до обробника (обробників). Це може бути, наприклад, споживач черги повідомлень або точка входу API.
Handler (обробник):
: Відповідає за обробку повідомлень із застосуванням бізнес-логіки, що стосується цих повідомлень. Обробники викликаються middleware HandleMessageMiddleware.
Middleware: : Middleware має доступ до повідомлення та його обгортки (конверта, envelope) під час його проходження через шину. Буквально «програмне забезпечення посередині» — воно не стосується основних завдань (бізнес-логіки) застосунку. Натомість це наскрізні аспекти, застосовні в межах усього застосунку і такі, що впливають на всю шину повідомлень. Наприклад: логування, валідація повідомлення, старт транзакції, ... Вони також відповідають за виклик наступного middleware в ланцюжку, а отже можуть змінювати конверт, додаючи до нього марки (stamps) або навіть замінюючи його, а також переривати ланцюжок middleware. Middleware викликається як тоді, коли повідомлення початково відправляється, так і пізніше, коли повідомлення отримано з транспорту.
Envelope (конверт): : Специфічна для Messenger концепція, що дає повну гнучкість усередині шини повідомлень, обгортаючи повідомлення й дозволяючи додавати всередину корисну інформацію через марки конверта (envelope stamps).
Envelope Stamps (марки конверта): : Частинка інформації, яку потрібно прикріпити до повідомлення: контекст серіалізатора для транспорту, маркери, що ідентифікують отримане повідомлення, або будь-які метадані, які може використовувати ваш middleware чи транспортний рівень.
Створення повідомлення та обробника
Messenger побудований довкола двох різних класів, які ви створюєте: (1) класу повідомлення, що містить дані, і (2) класу обробника (обробників), який буде викликаний, коли це повідомлення буде відправлено. Клас обробника читає клас повідомлення і виконує одне або кілька завдань.
До класу повідомлення немає особливих вимог, окрім того, що він має бути серіалізовним:
// src/Message/SmsNotification.php
namespace App\Message;
class SmsNotification
{
public function __construct(
private string $content,
) {
}
public function getContent(): string
{
return $this->content;
}
}
Обробник повідомлення — це PHP callable; рекомендований спосіб його створення — створити клас, що має атрибут Symfony\Component\Messenger\Attribute\AsMessageHandler і метод __invoke() із тип-хінтом класу повідомлення (або інтерфейсу повідомлення):
// src/MessageHandler/SmsNotificationHandler.php
namespace App\MessageHandler;
use App\Message\SmsNotification;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;
#[AsMessageHandler]
class SmsNotificationHandler
{
public function __invoke(SmsNotification $message)
{
// ... виконати якусь роботу — наприклад, надіслати SMS-повідомлення!
}
}
Порада. Ви також можете використовувати атрибут
#[AsMessageHandler]на окремих методах класу. Ви можете застосувати атрибут до будь-якої кількості методів в одному класі, що дозволяє згрупувати обробку кількох повʼязаних типів повідомлень.
Завдяки автоконфігурації та тип-хінту SmsNotification Symfony знає, що цей обробник має бути викликаний, коли відправляється повідомлення SmsNotification. У більшості випадків це все, що потрібно зробити. Але ви також можете налаштувати обробники повідомлень вручну. Щоб побачити всі налаштовані обробники, виконайте:
php bin/console debug:messenger
Відправлення повідомлення
Ви готові! Щоб відправити повідомлення (і викликати обробник), впровадьте сервіс messenger.default_bus (через MessageBusInterface), наприклад у контролері. Коли ви використовуєте Messenger у довільному PHP-застосунку, створіть шину повідомлень самостійно і зареєструйте в ній обробник(и) (у такому разі атрибут #[AsMessageHandler] не потрібен):
// src/Controller/DefaultController.php
namespace App\Controller;
use App\Message\SmsNotification;
use Symfony\Bundle\FrameworkBundle\Controller\AbstractController;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\Messenger\MessageBusInterface;
class DefaultController extends AbstractController
{
public function index(MessageBusInterface $bus): Response
{
// призведе до виклику SmsNotificationHandler
$bus->dispatch(new SmsNotification('Look! I created a message!'));
// ...
}
}
use App\Message\SmsNotification;
use App\MessageHandler\SmsNotificationHandler;
use Symfony\Component\Messenger\Handler\HandlersLocator;
use Symfony\Component\Messenger\MessageBus;
use Symfony\Component\Messenger\Middleware\HandleMessageMiddleware;
$bus = new MessageBus([
new HandleMessageMiddleware(new HandlersLocator([
SmsNotification::class => [new SmsNotificationHandler()],
])),
]);
// призведе до виклику SmsNotificationHandler
$bus->dispatch(new SmsNotification('Look! I created a message!'));
Транспорти: асинхронні повідомлення та повідомлення в черзі
За замовчуванням повідомлення обробляються одразу після відправлення. Якщо ви хочете обробляти повідомлення асинхронно, можете налаштувати транспорт. Транспорт здатний надсилати повідомлення (наприклад, до системи черг), а потім отримувати їх через обробник черги (worker). Messenger підтримує кілька транспортів.
Примітка. Якщо ви хочете використати транспорт, який не підтримується, погляньте на
транспорт Enqueue(Enqueue's transport), що підтримує такі сервіси, як Kafka та Google Pub/Sub.
Транспорт реєструється за допомогою «DSN». Завдяки рецепту Flex для Messenger у вашому файлі .env уже є кілька прикладів.
# MESSENGER_TRANSPORT_DSN=amqp://guest:guest@localhost:5672/%2f/messages
# MESSENGER_TRANSPORT_DSN=doctrine://default
# MESSENGER_TRANSPORT_DSN=redis://localhost:6379/messages
Розкоментуйте той транспорт, який вам потрібен (або задайте його в .env.local). Докладніше див. розділ про конфігурацію транспортів.
Далі, у config/packages/messenger.yaml, визначмо транспорт з назвою async, який використовує цю конфігурацію:
# config/packages/messenger.yaml
framework:
messenger:
transports:
async: "%env(MESSENGER_TRANSPORT_DSN)%"
# або в розгорнутому вигляді, щоб налаштувати більше опцій
#async:
# dsn: "%env(MESSENGER_TRANSPORT_DSN)%"
# options: []
// config/packages/messenger.php
namespace Symfony\Component\DependencyInjection\Loader\Configurator;
return App::config([
'framework' => [
'messenger' => [
'transports' => [
'async' => env('MESSENGER_TRANSPORT_DSN'),
// або в розгорнутому вигляді, щоб налаштувати більше опцій
// 'async' => [
// 'dsn' => env('MESSENGER_TRANSPORT_DSN'),
// 'options' => [],
// ],
],
],
],
]);
Маршрутизація повідомлень до транспорту
Тепер, коли транспорт налаштовано, замість негайної обробки повідомлення ви можете налаштувати їх надсилання до транспорту:
// src/Message/SmsNotification.php
namespace App\Message;
use Symfony\Component\Messenger\Attribute\AsMessage;
#[AsMessage('async')]
class SmsNotification
{
// ...
}
# config/packages/messenger.yaml
framework:
messenger:
transports:
async: "%env(MESSENGER_TRANSPORT_DSN)%"
routing:
# async — це будь-яка назва, яку ви дали своєму транспорту вище
'App\Message\SmsNotification': async
// config/packages/messenger.php
namespace Symfony\Component\DependencyInjection\Loader\Configurator;
use App\Message\SmsNotification;
return App::config([
'framework' => [
'messenger' => [
'routing' => [
// async — це будь-яка назва, яку ви дали своєму транспорту вище
SmsNotification::class => 'async',
],
],
],
]);
Завдяки цьому App\Message\SmsNotification буде надіслано до транспорту async, і його обробник(и) не буде викликано одразу. Будь-які повідомлення, що не збігаються з правилами в routing, і надалі оброблятимуться негайно, тобто синхронно.
Примітка. Якщо ви налаштовуєте маршрутизацію одночасно у файлах конфігурації YAML/PHP і через PHP-атрибути, конфігурація завжди має пріоритет над атрибутом класу. Така поведінка дозволяє перевизначати маршрутизацію окремо для кожного середовища.
Примітка. Налаштовуючи маршрутизацію в окремих файлах YAML/PHP, ви можете використати частковий простір імен PHP, як-от
'App\Message\*', щоб охопити всі повідомлення у відповідному просторі імен. Єдина вимога — символ підстановки'*'має стояти в кінці простору імен.Ви можете використати
'*'як клас повідомлення. Це працюватиме як типове правило маршрутизації для будь-якого повідомлення, що не збіглося з правилами вrouting. Це корисно, щоб гарантувати, що жодне повідомлення не оброблятиметься синхронно за замовчуванням.Єдиний недолік у тому, що
'*'застосується також до листів, надісланих Symfony Mailer (який використовуєSendEmailMessage, коли Messenger доступний). Це може спричинити проблеми, якщо ваші листи не серіалізовні (наприклад, якщо вони містять вкладення у вигляді PHP-ресурсів/потоків).
Ви також можете маршрутизувати класи за їхнім батьківським класом або інтерфейсом. Або надсилати повідомлення до кількох транспортів:
// src/Message/SmsNotification.php
namespace App\Message;
use Symfony\Component\Messenger\Attribute\AsMessage;
#[AsMessage(['async', 'audit'])]
class SmsNotification
{
// ...
}
// якщо вам так зручніше, ви також можете застосувати кілька атрибутів до класу повідомлення
#[AsMessage('async')]
#[AsMessage('audit')]
class SmsNotification
{
// ...
}
# config/packages/messenger.yaml
framework:
messenger:
routing:
# маршрутизувати всі повідомлення, що успадковують цей приклад базового класу або інтерфейсу
'App\Message\AbstractAsyncMessage': async
'App\Message\AsyncMessageInterface': async
'My\Message\ToBeSentToTwoSenders': [async, audit]
// config/packages/messenger.php
namespace Symfony\Component\DependencyInjection\Loader\Configurator;
use App\Message\AbstractAsyncMessage;
use App\Message\AsyncMessageInterface;
use My\Message\ToBeSentToTwoSenders;
return App::config([
'framework' => [
'messenger' => [
'routing' => [
// маршрутизувати всі повідомлення, що успадковують цей приклад базового класу або інтерфейсу
AbstractAsyncMessage::class => 'async',
AsyncMessageInterface::class => 'async',
ToBeSentToTwoSenders::class => ['async', 'audit'],
],
],
],
]);
Примітка. Якщо ви налаштуєте маршрутизацію і для дочірнього, і для батьківського класу, буде використано обидва правила. Наприклад, якщо у вас є обʼєкт
SmsNotification, що успадковуєNotification, буде використано маршрутизацію і дляNotification, і дляSmsNotification.
Порада. Ви можете визначати й перевизначати транспорт, який використовує повідомлення, під час виконання за допомогою
Symfony\Component\Messenger\Stamp\TransportNamesStampна конверті повідомлення. Ця марка приймає масив назв транспортів як єдиний аргумент. Докладніше про марки див. «Конверти та марки» (Envelopes & Stamps).
Сутності Doctrine у повідомленнях
Якщо вам потрібно передати сутність Doctrine у повідомленні, краще передати первинний ключ сутності (або будь-яку релевантну інформацію, яка справді потрібна обробнику, наприклад email тощо) замість обʼєкта (інакше ви можете побачити помилки, повʼязані з Entity Manager):
// src/Message/NewUserWelcomeEmail.php
namespace App\Message;
class NewUserWelcomeEmail
{
public function __construct(
private int $userId,
) {
}
public function getUserId(): int
{
return $this->userId;
}
}
Потім у своєму обробнику ви можете запитати свіжий обʼєкт:
// src/MessageHandler/NewUserWelcomeEmailHandler.php
namespace App\MessageHandler;
use App\Message\NewUserWelcomeEmail;
use App\Repository\UserRepository;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;
#[AsMessageHandler]
class NewUserWelcomeEmailHandler
{
public function __construct(
private UserRepository $userRepository,
) {
}
public function __invoke(NewUserWelcomeEmail $welcomeEmail): void
{
$user = $this->userRepository->find($welcomeEmail->getUserId());
// ... надіслати лист!
}
}
Це гарантує, що сутність містить свіжі дані.
Версіонування класів повідомлень
Клас повідомлення визначає контракт між кодом, який відправляє повідомлення, і обробником черги, який його обробляє. Оскільки Messenger обробляє повідомлення асинхронно, деякі повідомлення можуть усе ще очікувати в черзі, коли ви розгортаєте нову версію застосунку. Якщо ви зміните клас повідомлення, ті старіші повідомлення можуть більше не десеріалізуватися коректно, що може призвести до збоїв, коли обробники черги спробують їх обробити.
Для незначних змін зберігайте зворотну сумісність, роблячи нові аргументи конструктора необовʼязковими та надаючи розумне значення за замовчуванням:
final class SendInvoice
{
public function __construct(
public readonly int $orderId,
public readonly ?string $locale = null, // додано пізніше
) {
}
}
За такого підходу старіші повідомлення, що не містять нового аргументу $locale, і надалі можуть коректно десеріалізуватися.
Ще одна зміна, що ламає сумісність, — це видалення властивостей із класу повідомлення. Повідомлення, вже збережені в черзі, можуть усе ще містити ці властивості, коли їх десеріалізують новіші обробники черги. Починаючи з PHP 8.2, це може спричинити попередження про застарілість, оскільки динамічні властивості є застарілими.
Якщо вам необхідно видалити властивість, розгляньте один із таких підходів:
- Тимчасово залишити властивість (наприклад, як
public ?type $property = null), доки всі старі повідомлення не буде оброблено. - Додати до класу атрибут
#[\AllowDynamicProperties], щоб дозволити старішим серіалізованим повідомленням задавати властивості, яких більше не існує. - Реалізувати власну логіку серіалізації, щоб контролювати, як повідомлення серіалізується й десеріалізується.
Якщо зміна змінює сенс повідомлення, а не просто розширює його, створіть нову версію класу повідомлення замість того, щоб змінювати наявний:
// Зберігайте SendInvoice, доки не буде оброблено всі повідомлення в черзі, що його використовують
final class SendInvoiceV2
{
public function __construct(
public readonly int $orderId,
public readonly string $locale,
public readonly string $templateId,
) {
}
}
Під час розгортання обидві версії можуть тимчасово співіснувати. Спочатку розгорніть новий клас повідомлення SendInvoiceV2 та його обробник, зберігаючи старі. Після того як чергу буде повністю спорожнено, видаліть старий клас і обробник у наступному розгортанні.
Занадто раннє видалення старого класу може призвести до збоїв обробників черги, коли вони намагатимуться десеріалізувати повідомлення, відправлені до розгортання.
Синхронна обробка повідомлень
Якщо повідомлення не збігається з жодним правилом маршрутизації, воно не буде надіслане до жодного транспорту і буде оброблене негайно. У деяких випадках (наприклад, коли привʼязуєте обробники до різних транспортів) простіше або гнучкіше зробити це явно: створити транспорт sync і «надсилати» повідомлення туди, щоб вони оброблялися негайно:
# config/packages/messenger.yaml
framework:
messenger:
transports:
# ... інші транспорти
sync: 'sync://'
routing:
App\Message\SmsNotification: sync
// config/packages/messenger.php
namespace Symfony\Component\DependencyInjection\Loader\Configurator;
use App\Message\SmsNotification;
return App::config([
'framework' => [
'messenger' => [
'transports' => [
'sync' => 'sync://',
],
'routing' => [
SmsNotification::class => 'sync',
],
],
],
]);
Створення власного транспорту
Ви також можете створити власний транспорт, якщо вам потрібно надсилати чи отримувати повідомлення від чогось, що не підтримується. Див. /messenger/custom-transport.
Споживання повідомлень (запуск обробника черги)
Щойно ваші повідомлення промаршрутизовано, у більшості випадків вам потрібно буде їх «спожити». Це можна зробити командою messenger:consume:
php bin/console messenger:consume async
# використайте -vv, щоб побачити подробиці того, що відбувається
php bin/console messenger:consume async -vv
# використайте регулярний вираз, щоб охопити кілька транспортів одразу
# (наприклад, усі транспорти, що починаються з "scheduler_")
php bin/console messenger:consume scheduler_.*
# поєднуйте регулярні вирази та явні назви транспортів
php bin/console messenger:consume high_priority.* low_priority async
Перший аргумент — це назва отримувача (або id сервісу, якщо ви маршрутизували до власного сервісу). За замовчуванням команда працюватиме вічно: шукатиме нові повідомлення у вашому транспорті й оброблятиме їх. Ця команда й називається вашим «обробником черги» (worker).
Ви також можете використовувати регулярні вирази як назву отримувача, щоб охопити кілька транспортів одразу. Це корисно, коли у вас є кілька транспортів зі схожими назвами, наприклад транспорти, згруповані за призначенням або пріоритетом.
Коли регулярний вираз відповідає кільком транспортам, вони споживаються в тому порядку, у якому визначені у вашій конфігурації. Якщо ви вкажете кілька назв отримувачів або регулярних виразів, вони обробляються в тому порядку, у якому ви передаєте їх команді.
Примітка. Зіставлення з регулярними виразами працює лише тоді, коли команда не запущена в інтерактивному режимі.
Нове у версії 8.1. Підтримку регулярних виразів як назви отримувача було додано в Symfony 8.1.
Якщо ви хочете споживати повідомлення з усіх доступних отримувачів, можете використати команду з опцією --all:
php bin/console messenger:consume --all
Використовуючи --all, ви можете виключити окремих отримувачів за допомогою опції --exclude-receivers (скорочення -eq):
php bin/console messenger:consume --all --exclude-receivers=async_priority_low --exclude-receivers=failed
Примітка. Опцію
--exclude-receiversможна використовувати лише разом із--all. Крім того, ви не можете виключити всіх отримувачів.
Повідомлення, обробка яких триває довго, можуть бути передчасно доставлені повторно, оскільки деякі транспорти вважають, що непідтверджене повідомлення втрачене. Щоб запобігти цій проблемі, використовуйте опцію команди --keepalive, щоб задати інтервал (у секундах; значення за замовчуванням = 5), з яким повідомлення позначається як «у процесі». Це запобігає повторній доставці повідомлення, доки обробник черги не завершить його обробку:
php bin/console messenger:consume --keepalive
Примітка. Ця опція доступна лише для таких транспортів: Beanstalkd, AmazonSQS, Doctrine і Redis.
За замовчуванням обробник черги забирає з транспорту одне повідомлення за ітерацію. Використовуйте опцію --fetch-size, щоб забирати кілька повідомлень за ітерацію, зменшуючи кількість звернень до транспорту:
php bin/console messenger:consume async --fetch-size=8
Це особливо корисно з транспортами, які підтримують отримання кількох повідомлень за один виклик, як-от Amazon SQS, Redis, AMQP і Doctrine.
Нове у версії 8.1. Опцію
--fetch-sizeбуло додано в Symfony 8.1.
Порада. У середовищі розробки, якщо ви користуєтеся інструментом Symfony CLI, ви можете налаштувати автоматичний запуск обробників черги разом із вебсервером. Докладніше — у документації «Symfony CLI Workers».
Порада. Щоб коректно зупинити обробник черги, викиньте екземпляр
Symfony\Component\Messenger\Exception\StopWorkerException.
Розгортання у продакшн
У продакшні є кілька важливих речей, про які варто подумати:
Використовуйте менеджер процесів, як-от Supervisor або systemd, щоб ваші обробники черги працювали постійно
Вам потрібно, щоб один чи більше «обробників черги» працювали весь час. Для цього використовуйте систему контролю процесів, як-от Supervisor або systemd.
Не дозволяйте обробникам черги працювати вічно
Деякі сервіси (як-от EntityManager від Doctrine) з часом споживають дедалі більше памʼяті. Тож замість того, щоб дозволяти обробнику черги працювати вічно, використайте прапорець на кшталт messenger:consume --limit=10, щоб сказати обробнику обробити лише 10 повідомлень перед виходом (після чого менеджер процесів створить новий процес). Є також інші опції, як-от --memory-limit=128M і --time-limit=3600.
Зупинка обробників черги, що натрапляють на помилки
Якщо залежність обробника черги, наприклад ваш сервер бази даних, недоступна або досягнуто таймауту, ви можете спробувати додати логіку повторного підключення або просто завершувати обробник, якщо він отримує забагато помилок, — за допомогою опції --failure-limit команди messenger:consume.
Перезапускайте обробники черги під час розгортання
Щоразу під час розгортання перезапускайте всі процеси обробників черги, щоб вони підхопили новий код. Для цього виконайте messenger:stop-workers після того, як новий код опиниться на диску й кеш буде прогрітий, але до перемикання трафіку. Це встановлює прапорець у кеші, який каже кожному обробнику завершити поточне повідомлення й коректно вийти. Далі менеджер процесів перезапускає їх уже на новій кодовій базі:
# у вашому скрипті розгортання, після викладення коду й прогріву кешу:
php bin/console messenger:stop-workers
Команда внутрішньо використовує кеш застосунку (app cache). Якщо ваш застосунок працює на кількох хостах, налаштуйте кеш застосунку на використання спільного адаптера (наприклад, Redis), щоб усі вебпроцеси та процеси обробників черги використовували один і той самий кеш.
Примітка. У середовищі Kubernetes послідовний перезапуск (rolling restart)
Deploymentобробників черги дає той самий результат, але лише якщоterminationGracePeriodSecondsдостатньо великий, щоб найдовший обробник встиг завершитися до заміни пода.SIGKILLне дає обробникам черги шансу завершити поточне повідомлення, через що обробник може лишитися посеред виконання.
Використовуйте той самий кеш між розгортаннями
Якщо ваша стратегія розгортання передбачає створення нових цільових директорій, вам слід задати значення для опції конфігурації cache.prefix_seed, щоб використовувати той самий простір імен кешу між розгортаннями. Інакше пул cache.app використовуватиме значення параметра kernel.project_dir як основу для простору імен, що призводитиме до різних просторів імен щоразу під час нового розгортання.
Пріоритезовані транспорти
Використовуйте окремі транспорти для типів повідомлень із різними вимогами до затримки, різними режимами відмов чи вікнами повторних спроб. Коли кілька типів повідомлень поділяють один транспорт, повільний чи проблемний обробник одного типу може затримувати всі інші в тій самій черзі.
Наприклад, якщо повідомлення синхронізації каталогу й підтвердження платежів маршрутизуються до одного транспорту, повільний обробник товарного фіду може затримати підтвердження платежів для клієнтів, які проходять оформлення замовлення.
Розгляньте призначення кожному транспорту власного обробника черги, щоб збої чи сповільнення в одному потоці повідомлень не впливали на інші:
# config/packages/messenger.yaml
framework:
messenger:
transports:
async_priority_high:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
# queue_name є специфічним для транспорту doctrine
queue_name: high
# для AMQP надсилайте до окремого exchange, потім до черги
#exchange:
# name: high
#queues:
# messages_high: ~
# для redis спробуйте "group"
async_priority_low:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
queue_name: low
routing:
'App\Message\SmsNotification': async_priority_low
'App\Message\NewUserWelcomeEmail': async_priority_high
// config/packages/messenger.php
namespace Symfony\Component\DependencyInjection\Loader\Configurator;
use App\Message\NewUserWelcomeEmail;
use App\Message\SmsNotification;
return App::config([
'framework' => [
'messenger' => [
'transports' => [
'async_priority_high' => [
'dsn' => env('MESSENGER_TRANSPORT_DSN'),
'options' => [
// queue_name є специфічним для транспорту doctrine
'queue_name' => 'high',
// для AMQP надсилайте до окремого exchange, потім до черги
// 'exchange' => [
// 'name' => 'high',
// ],
// 'queues' => [
// 'messages_high' => null,
// ],
// для redis спробуйте "group"
],
],
'async_priority_low' => [
'dsn' => env('MESSENGER_TRANSPORT_DSN'),
'options' => [
'queue_name' => 'low',
],
],
],
'routing' => [
SmsNotification::class => 'async_priority_low',
NewUserWelcomeEmail::class => 'async_priority_high',
],
],
],
]);
Далі ви можете запускати окремі обробники черги для кожного транспорту або вказати одному обробнику обробляти повідомлення в порядку пріоритету:
php bin/console messenger:consume async_priority_high async_priority_low
Обробник черги завжди спершу шукатиме повідомлення, що очікують в async_priority_high. Якщо їх немає, тоді він споживатиме повідомлення з async_priority_low.
Пріоритезовані повідомлення
Нове у версії 8.1. Підтримку пріоритезації повідомлень було додано в Symfony 8.1.
За замовчуванням Messenger використовує порядок «перший прийшов — перший вийшов», тож повідомлення отримуються в тому самому порядку, у якому були надіслані, за винятком відкладених повідомлень.
Пріоритезовані транспорти дають базову підтримку пріоритетів, але лише між різними типами повідомлень. Використовуючи транспорти AMQP або Beanstalkd, ви можете пріоритезувати повідомлення одного типу за допомогою марки під назвою PriorityStamp:
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Stamp\PriorityStamp;
$bus->dispatch(
(new Envelope($message))->with(new PriorityStamp(255))
);
Порада. Підтримувані пріоритети — у діапазоні від
0(найнижчий) до255(найвищий).
PriorityStamp можна поєднувати з DelayStamp. Коли затримка спливає, повідомлення доставляється раніше за будь-які повідомлення з нижчим пріоритетом, що вже перебувають у черзі:
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Stamp\DelayStamp;
use Symfony\Component\Messenger\Stamp\PriorityStamp;
$bus->dispatch(
(new Envelope($message))->with(
new PriorityStamp(255),
new DelayStamp(5000)
)
);
Під час використання транспорту Beanstalkd додаткова конфігурація не потрібна. Beanstalkd внутрішньо використовує обернену шкалу пріоритетів (0 = найвищий, 2^32 - 1 = найнижчий), але Messenger виконує перетворення автоматично.
Під час використання транспорту AMQP черги з пріоритетами мають бути явно увімкнені в конфігурації:
# config/packages/messenger.yaml
framework:
messenger:
transports:
async:
dsn: "%env(MESSENGER_TRANSPORT_DSN)%"
options:
queues:
messenger:
arguments:
x-max-priority: 255
// config/packages/messenger.php
use Symfony\Config\FrameworkConfig;
return static function (FrameworkConfig $framework): void {
$framework->messenger()
->transport('async')
->dsn('%env(MESSENGER_TRANSPORT_DSN)%')
->options([
'queues' => [
'messenger' => [
'arguments' => ['x-max-priority' => 255],
],
],
]);
};
Увага!
x-max-priorityне можна змінити на наявній черзі RabbitMQ. Messenger не зможе виконати автоматичне налаштування, і пріоритети не працюватимуть. Натомість створіть нову чергу.
Примітка. RabbitMQ рекомендує використовувати
не більше 10 рівнів пріоритету(no more than 10 priority levels). Наприклад:255для високого,127для середнього і0для низького. Більша кількість рівнів може вплинути на продуктивність.
Обмеження споживання конкретними чергами
Деякі транспорти (зокрема AMQP) мають поняття exchange і черг. Транспорт Symfony завжди привʼязаний до exchange. За замовчуванням обробник черги споживає з усіх черг, приєднаних до exchange вказаного транспорту. Проте бувають випадки, коли потрібно, щоб обробник споживав лише з певних черг.
Ви можете обмежити обробник черги, щоб він обробляв повідомлення лише з конкретної черги (черг):
php bin/console messenger:consume my_transport --queues=fasttrack
# ви можете передати опцію --queues більш ніж один раз, щоб обробляти кілька черг
php bin/console messenger:consume my_transport --queues=fasttrack1 --queues=fasttrack2
Примітка. Щоб можна було використовувати опцію
queues, отримувач має реалізовуватиSymfony\Component\Messenger\Transport\Receiver\QueueReceiverInterface.
Перевірка кількості повідомлень у черзі для кожного транспорту
Виконайте команду messenger:stats, щоб дізнатися, скільки повідомлень перебуває в «чергах» деяких або всіх транспортів:
# показує кількість повідомлень у черзі в усіх транспортах
php bin/console messenger:stats
# показує статистику лише для деяких транспортів
php bin/console messenger:stats my_transport_name other_transport_name
# ви також можете вивести статистику у форматі JSON
php bin/console messenger:stats --format=json
php bin/console messenger:stats my_transport_name other_transport_name --format=json
Примітка. Щоб ця команда працювала, отримувач налаштованого транспорту має реалізовувати
Symfony\Component\Messenger\Transport\Receiver\MessageCountAwareInterface.
Конфігурація Supervisor
Supervisor — чудовий інструмент, щоб гарантувати, що ваш процес (процеси) обробника черги завжди працює (навіть якщо він завершується через збій, досягнення ліміту повідомлень або завдяки messenger:stop-workers). Ви можете встановити його, наприклад, в Ubuntu через:
sudo apt-get install supervisor
Файли конфігурації Supervisor зазвичай розташовані в директорії /etc/supervisor/conf.d. Наприклад, ви можете створити там новий файл messenger-worker.conf, щоб переконатися, що 2 екземпляри messenger:consume працюють увесь час:
;/etc/supervisor/conf.d/messenger-worker.conf
[program:messenger-consume]
command=php /path/to/your/app/bin/console messenger:consume async --time-limit=3600
user=ubuntu
numprocs=2
startsecs=0
autostart=true
autorestart=true
startretries=10
process_name=%(program_name)s_%(process_num)02d
Змініть аргумент async на назву вашого транспорту (або транспортів), а user — на Unix-користувача на вашому сервері.
Увага! Під час розгортання щось може бути недоступним (наприклад, база даних), через що споживач не зможе запуститися. У такій ситуації Supervisor спробує перезапустити команду
startretriesразів. Обовʼязково змініть цю настройку, щоб команда не потрапила у стан FATAL, з якого вона більше ніколи не перезапуститься.З кожним перезапуском Supervisor збільшує затримку на 1 секунду. Наприклад, якщо значення дорівнює
10, він чекатиме 1 с, 2 с, 3 с тощо. Це дає сервісу загалом 55 секунд, щоб знову стати доступним. Збільште настройкуstartretries, щоб покрити максимально очікуваний простій.
Якщо ви використовуєте транспорт Redis, зауважте, що кожному обробнику черги потрібне унікальне імʼя споживача, щоб уникнути обробки одного й того самого повідомлення кількома обробниками. Один зі способів цього досягти — задати змінну середовища у файлі конфігурації Supervisor, на яку потім можна посилатися в messenger.yaml (див. розділ про Redis нижче):
environment=MESSENGER_CONSUMER_NAME=%(program_name)s_%(process_num)02d
Далі скажіть Supervisor прочитати вашу конфігурацію і запустити ваші обробники черги:
sudo supervisorctl reread
sudo supervisorctl update
sudo supervisorctl start messenger-consume:*
# Якщо ви розгортаєте оновлення свого коду, не забудьте перезапустити обробники черги,
# щоб запустити новий код
sudo supervisorctl restart messenger-consume:*
Докладніше — у документації Supervisor (Supervisor docs).
Коректне завершення роботи
Якщо ви встановите у своєму проєкті PHP-розширення PCNTL, обробники черги оброблятимуть POSIX-сигнали SIGTERM або SIGINT, щоб завершити обробку поточного повідомлення перед припиненням роботи.
Проте ви можете віддавати перевагу іншим POSIX-сигналам для коректного завершення. Ви можете перевизначити типові, задавши опцію конфігурації framework.messenger.stop_worker_on_signals:
# config/packages/messenger.yaml
framework:
messenger:
stop_worker_on_signals:
- SIGTERM
- SIGINT
- SIGUSR1
# ...
<!-- config/packages/messenger.xml -->
<?xml version="1.0" encoding="UTF-8" ?>
<container xmlns="http://symfony.com/schema/dic/services"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:framework="http://symfony.com/schema/dic/symfony"
xsi:schemaLocation="http://symfony.com/schema/dic/services
https://symfony.com/schema/dic/services/services-1.0.xsd
http://symfony.com/schema/dic/symfony
https://symfony.com/schema/dic/symfony/symfony-1.0.xsd">
<framework:config>
<framework:messenger>
<!-- ... -->
<framework:stop-worker-on-signals signal="SIGTERM"/>
<framework:stop-worker-on-signals signal="SIGINT"/>
<framework:stop-worker-on-signals signal="SIGUSR1"/>
</framework:messenger>
</framework:config>
</container>
// config/packages/messenger.php
namespace Symfony\Component\DependencyInjection\Loader\Configurator;
return App::config([
'framework' => [
'messenger' => [
'stop_worker_on_signals' => [
'SIGTERM',
'SIGINT',
'SIGUSR1',
],
],
],
]);
У деяких випадках сигнал SIGTERM надсилає сам Supervisor (наприклад, під час зупинки Docker-контейнера, у якого Supervisor є entrypoint). У таких випадках вам потрібно додати ключ stopwaitsecs до конфігурації програми (зі значенням бажаного пільгового періоду в секундах), щоб виконати коректне завершення роботи:
[program:x]
stopwaitsecs=20
Конфігурація systemd
Хоча Supervisor — чудовий інструмент, він має недолік: щоб його запустити, потрібен доступ до системи. Systemd став стандартом у більшості дистрибутивів Linux і має хорошу альтернативу — користувацькі сервіси (user services).
Файли конфігурації користувацьких сервісів systemd зазвичай розташовані в директорії ~/.config/systemd/user. Наприклад, ви можете створити новий файл messenger-worker.service. Або файл [email protected], якщо хочете, щоб одночасно працювало більше екземплярів:
[Unit]
Description=Symfony messenger-consume %i
[Service]
ExecStart=php /path/to/your/app/bin/console messenger:consume async --time-limit=3600
# для Redis задайте власне імʼя споживача для кожного екземпляра
Environment="MESSENGER_CONSUMER_NAME=symfony-%n-%i"
Restart=always
RestartSec=30
[Install]
WantedBy=default.target
Тепер скажіть systemd увімкнути й запустити один обробник черги:
systemctl --user enable [email protected]
systemctl --user start [email protected]
# щоб увімкнути й запустити 20 обробників черги
systemctl --user enable messenger-worker@{1..20}.service
systemctl --user start messenger-worker@{1..20}.service
Якщо ви зміните файл конфігурації сервісу, вам потрібно перезавантажити демон:
systemctl --user daemon-reload
Щоб перезапустити всіх ваших споживачів:
systemctl --user restart messenger-consume@*.service
Користувацький екземпляр systemd запускається лише після першого входу відповідного користувача в систему. Натомість споживачам часто потрібно запускатися під час завантаження системи. Увімкніть lingering для користувача, щоб активувати таку поведінку:
loginctl enable-linger <your-username>
Логами керує journald, і з ними можна працювати за допомогою команди journalctl:
# стежити за логами споживача nr 11
journalctl -f --user-unit [email protected]
# стежити за логами всіх споживачів
journalctl -f --user-unit messenger-consume@*
# стежити за всіма логами ваших користувацьких сервісів
journalctl -f _UID=$UID
Докладніше — у документації systemd (systemd docs).
Примітка. Вам потрібні або підвищені привілеї для команди
journalctl, або додайте свого користувача до групи systemd-journal:sudo usermod -a -G systemd-journal <your-username>
Обробник черги без стану
PHP спроєктовано як мову без стану; між різними запитами немає спільних ресурсів. У HTTP-контексті PHP очищає все після надсилання відповіді, тож ви можете вирішити не перейматися сервісами, які можуть спричиняти витік памʼяті.
З іншого боку, для обробників черги звично обробляти повідомлення послідовно в довготривалих CLI-процесах, які не завершуються після обробки одного повідомлення. Остерігайтеся станів сервісів, щоб запобігти витоку інформації та/або памʼяті, оскільки Symfony впроваджуватиме той самий екземпляр сервісу в усіх повідомленнях, зберігаючи внутрішній стан сервісів.
Проте певні сервіси Symfony, як-от Monolog fingers crossed handler, «течуть» за задумом. Symfony надає можливість скидання сервісів (service reset), щоб розвʼязати цю проблему. Автоматично скидаючи контейнер між двома повідомленнями, Symfony шукає будь-які сервіси, що реалізують Symfony\Contracts\Service\ResetInterface (включно з вашими власними сервісами), і викликає їхній метод reset(), щоб вони могли очистити свій внутрішній стан.
Якщо сервіс не є безстановим і ви хочете скидати його властивості після кожного повідомлення, то цей сервіс має реалізовувати Symfony\Contracts\Service\ResetInterface, де ви можете скинути властивості в методі reset().
Якщо ви не хочете скидати контейнер, додайте опцію --no-reset під час запуску команди messenger:consume.
Замість того щоб повністю вимикати скидання, ви також можете налаштувати, щоб скидання відбувалося кожні N повідомлень, передавши число в опцію --no-reset. Це корисно, коли скидання після кожного повідомлення надто дороге, але повна відсутність скидання призводить до витоків памʼяті чи застарілого стану:
# скидати сервіси кожні 100 повідомлень замість після кожного
php bin/console messenger:consume async --no-reset=100
Нове у версії 8.1. Можливість передавати число в опцію
--no-resetбуло додано в Symfony 8.1.
Власна стратегія виконання повідомлень
Нове у версії 8.1.
MessageExecutionStrategyInterfaceбуло додано в Symfony 8.1.
За замовчуванням обробник черги обробляє повідомлення синхронно, використовуючи Symfony\Component\Messenger\Execution\SyncMessageExecutionStrategy.
Неперекладеними лишилися підрозділи, що йдуть далі в оригіналі: продовження «Custom Message Execution Strategy», «Retries & Failures» (Retries: Handling Transient Failures, Avoiding Retrying Non-Recoverable Errors, Forcing Retries, Saving & Retrying Failed Messages), «Transports: Configuration» (AMQP Transport, Doctrine Transport, Redis Transport, In Memory Transport, Amazon SQS, Beanstalkd Transport, Serializing Messages), «Customizing Handlers» (Manually Configuring Handlers, Handler Subscriber & Options, Binding Handlers to Different Transports, Extending Messenger), «Middleware» (Adding your own Middleware, Middleware for Doctrine), «Messenger Events», «Multiple Buses, Command & Event Buses», «Learn more».
Перекладаємо з офіційної документації, розділ за розділом, і не ховаємо недоперекладене. Помітили неточність у терміні чи реченні: напишіть, виправимо.