W architekturze opartej na zdarzeniach jednym z kluczowych elementów jest możliwość asynchronicznego przetwarzania zadań. RabbitMQ i Symfony Messenger są świetnymi narzędziami do obsługi komunikatów, a w tym artykule skupimy się na konfiguracji konsumentów, implementacji procesów przetwarzania i tworzeniu dedykowanych workerów do przetwarzania różnych rodzajów komunikatów. Omówimy krok po kroku, jak skonfigurować kolejki, konsumentów i jak zarządzać workerami.
Architektura Zdarzeń i Kolejki RabbitMQ
Kolejki w RabbitMQ to mechanizm umożliwiający przechowywanie komunikatów aż do momentu ich przetworzenia przez konsumenta. Dzięki temu możemy oddzielić moment wysłania komunikatu od jego przetwarzania, co daje nam elastyczność i niezależność między komponentami systemu.
Konsument (ang. consumer) jest odpowiedzialny za odbieranie komunikatów z kolejki i ich przetwarzanie. W Symfony Messenger konsument może być skonfigurowany tak, aby odbierać wiadomości i wykonywać konkretne zadania.
Krok 1: Instalacja Symfony Messenger i RabbitMQ
Pierwszym krokiem jest upewnienie się, że mamy zainstalowane odpowiednie komponenty. Symfony Messenger z RabbitMQ wymaga enqueue/amqp-bunny. Możemy zainstalować go przy pomocy Composer:
composer require symfony/messenger enqueue/amqp-bunny
Krok 2: Konfiguracja Konsumentów w Symfony Messenger
Po instalacji należy skonfigurować Messenger, aby mógł połączyć się z RabbitMQ. Poniżej przedstawiam przykładową konfigurację:
# config/packages/messenger.yaml
framework:
messenger:
transports:
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
exchange:
name: 'messages'
routing:
'App\Message\TaskMessage': async
- dsn: To identyfikator źródła danych, który mówi Symfony, jak połączyć się z RabbitMQ. Można go skonfigurować w
.env:
MESSENGER_TRANSPORT_DSN=amqp://guest:guest@localhost:5672/%2f/messages
routing: Określa, które klasy wiadomości mają być wysyłane do danej kolejki. Tutaj
TaskMessagebędzie kierowana doasync.
Krok 3: Tworzenie Klasy Wiadomości
W celu przetwarzania zadań musimy stworzyć klasę wiadomości, która będzie reprezentować konkretne zadanie do wykonania:
namespace App\Message;
class TaskMessage
{
private string $taskId;
public function __construct(string $taskId)
{
$this->taskId = $taskId;
}
public function getTaskId(): string
{
return $this->taskId;
}
}
Ta klasa TaskMessage zawiera jedno pole $taskId, które identyfikuje konkretne zadanie do wykonania.
Krok 4: Implementacja Handlera Wiadomości
Handler (konsument) jest odpowiedzialny za odbieranie komunikatów i przetwarzanie ich. Stworzymy klasę TaskMessageHandler:
namespace App\MessageHandler;
use App\Message\TaskMessage;
use Symfony\Component\Messenger\Handler\MessageHandlerInterface;
class TaskMessageHandler implements MessageHandlerInterface
{
public function __invoke(TaskMessage $message)
{
// Przykład przetwarzania zadania na podstawie $taskId
$taskId = $message->getTaskId();
// Załóżmy, że przetwarzamy jakieś zadanie (np. generowanie raportu)
echo "Przetwarzanie zadania o ID: " . $taskId . "\n";
// Logika przetwarzania zadania
// ... (np. wysyłka e-maila, aktualizacja bazy danych, itp.)
}
}
Handler TaskMessageHandler implementuje interfejs MessageHandlerInterface, co oznacza, że jest automatycznie wywoływany, gdy w kolejce pojawia się komunikat TaskMessage. W metodzie __invoke() możemy realizować dowolną logikę przetwarzania.
Krok 5: Uruchamianie Konsumenta
Aby konsument mógł odbierać wiadomości z kolejki, musimy uruchomić tzw. workera. Robimy to za pomocą polecenia:
php bin/console messenger:consume async
asyncto nazwa transportu (kolejki), której worker ma słuchać.- Konsument będzie działał w tle, przetwarzając kolejne wiadomości z kolejki, gdy się pojawią.
Krok 6: Tworzenie Dedykowanych Workerów
Możemy również skonfigurować różnych workerów do przetwarzania różnych rodzajów wiadomości. Załóżmy, że mamy dodatkowy typ wiadomości EmailMessage, który ma być przetwarzany przez osobny worker.
Tworzenie Klasy Wiadomości EmailMessage
namespace App\Message;
class EmailMessage
{
private string $email;
public function __construct(string $email)
{
$this->email = $email;
}
public function getEmail(): string
{
return $this->email;
}
}
Tworzenie Handlera EmailMessageHandler
namespace App\MessageHandler;
use App\Message\EmailMessage;
use Symfony\Component\Messenger\Handler\MessageHandlerInterface;
class EmailMessageHandler implements MessageHandlerInterface
{
public function __invoke(EmailMessage $message)
{
// Przykład wysyłania e-maila
$email = $message->getEmail();
echo "Wysyłanie e-maila do: " . $email . "\n";
// Wyobraźmy sobie, że używamy jakiejś usługi do wysyłki e-maila
mail($email, 'Subject', 'Email content');
}
}
Aktualizacja Konfiguracji Messenger
Aby skonfigurować routing dla EmailMessage, dodajemy wpis w messenger.yaml:
framework:
messenger:
routing:
'App\Message\TaskMessage': async
'App\Message\EmailMessage': async
Teraz wiadomości EmailMessage i TaskMessage będą kierowane do tej samej kolejki async.
Uruchamianie Dedykowanych Workerów
Możemy uruchomić dedykowanych workerów do przetwarzania różnych rodzajów wiadomości:
Worker dla TaskMessage:
php bin/console messenger:consume async --limit=10 --time-limit=300 --no-reset
Ten worker przetworzy maksymalnie 10 wiadomości lub zatrzyma się po 5 minutach. Opcja --no-reset sprawia, że worker nie resetuje stanu, co przydatne jest w niektórych sytuacjach.
- Worker dla EmailMessage: Aby uruchomić innego workera, możemy dodać filtr do przetwarzania tylko
EmailMessage:
php bin/console messenger:consume async --receivers=App\Message\EmailMessage
Dzięki temu mamy elastyczność w zarządzaniu workerami, co pozwala na lepsze skalowanie systemu.
Krok 7: Przegląd Strategii Odbierania Wiadomości
Symfony Messenger daje możliwość definiowania strategii przetwarzania komunikatów:
Retry: Gdy przetwarzanie wiadomości się nie uda, możemy zdefiniować strategię ponowienia, np. ile razy i z jakim opóźnieniem próbować ponownie.
framework:
messenger:
transports:
async:
retry_strategy:
max_retries: 3
delay: 1000
multiplier: 2
max_delay: 10000
Ta konfiguracja oznacza, że po nieudanym przetwarzaniu system spróbuje jeszcze 3 razy, przy czym każde kolejne opóźnienie będzie rosnąć dwukrotnie.
Delay: Możemy także opóźniać przetwarzanie wiadomości, co jest przydatne w niektórych scenariuszach, np. czekanie na jakieś zewnętrzne dane.
Failure Transport: Jeśli przetwarzanie wiadomości się nie uda po maksymalnej liczbie prób, możemy przekierować ją do osobnej kolejki (np.
failed), aby przetworzyć ją ręcznie później.
Podsumowanie
RabbitMQ i Symfony Messenger to potężne narzędzia do tworzenia asynchronicznej architektury opartej na zdarzeniach. Dzięki kolejkom i konsumentom możemy rozdzielać złożone zadania i przetwarzać je niezależnie, co znacznie zwiększa skalowalność i odporność aplikacji.
W tym artykule pokazaliśmy, jak skonfigurować konsumentów, implementować przetwarzanie wiadomości i tworzyć dedykowanych workerów. Zrozumienie i efektywne wykorzystanie tych technik pozwala na tworzenie aplikacji, które są skalowalne, łatwe w utrzymaniu, a także odporne na obciążenia i awarie.
Jeśli jesteś zainteresowany architekturą event-driven i chcesz budować aplikacje gotowe na wyzwania współczesnego świata, warto zapoznać się z RabbitMQ i Symfony Messenger. Dzięki temu możesz uniezależnić komponenty systemu, sprawić, że przetwarzanie stanie się bardziej elastyczne i odporniejsze na niepowodzenia.