RabbitMQ i Symfony: Kolejki i konsumenci w architekturze zdarzeń

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 TaskMessage będzie kierowana do async.

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

				
			
  • async to 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:

  1. 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.

  1. 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.