Endpoint checkoutu, który zmniejsza obraz produktu, wysyła email z potwierdzeniem i aktualizuje silnik rekomendacji przed odpowiedzią 200 OK, jest tak szybki jak jego najwolniejszy krok. Jeśli którykolwiek z nich przekroczy limit czasu, cały request się wywala. Kolejka wiadomości pozwala endpointowi oddać tę pracę jako wiadomości i od razu odpowiedzieć, podczas gdy osobne workery odbierają je we własnym tempie.
Co faktycznie robi kolejka wiadomości
Kolejka wiadomości siedzi między dwiema częściami systemu, które nie muszą działać w tym samym momencie. Producer publikuje wiadomość opisującą pracę do wykonania. Kolejka ją przechowuje. Consumer ją odbiera, przetwarza i potwierdza. Te dwie strony nigdy nie rozmawiają ze sobą bezpośrednio.
Pomyśl o skrzynce pocztowej. Nadawca wrzuca list i odchodzi; odbiorca sprawdza skrzynkę, kiedy ma czas, a jeśli nie ma go godzinę, listy po prostu czekają.
Z tego wynikają trzy rzeczy, których bezpośrednie wywołanie funkcji dać nie może:
- Rozdzielenie. Endpoint checkoutu nie musi wiedzieć, jak działa zmniejszanie obrazów ani że to wolne. Publikuje
{"event": "order_placed", "order_id": 4821}i idzie dalej. - Buforowanie. Jeśli w tej samej sekundzie przychodzi 500 zamówień, kolejka wchłania ten skok. Workery przetwarzają je w tempie, jakie są w stanie utrzymać, zamiast blokować 500 requestów naraz.
- Niezależne skalowanie i awarie. Uruchom trzy workery do zmniejszania obrazów i jeden do maili, zrestartuj jeden bez dotykania drugiego, wdróż serwis checkoutu bez ponownego wdrażania jego workerów.
Uruchamianie RabbitMQ w Dockerze
RabbitMQ to najpopularniejszy uniwersalny message broker i sensowny domyślny wybór, jeśli nie masz jeszcze żadnej kolejki. Pobierz obraz management, żeby oprócz brokera dostać też webowy UI:
| |
Port 5672 to protokół AMQP, z którym łączy się twoja aplikacja. Port 15672 to UI do zarządzania: otwórz http://localhost:15672, zaloguj się jako app / changeme i na żywo obserwuj kolejki, częstotliwość wiadomości i połączenia. Domyślny login guest/guest działa tylko wewnątrz kontenera, więc ustaw własnego użytkownika do wszystkiego, z czym łączysz się z hosta.
Jeśli nigdy nie używałeś Dockera, Co to jest Docker i do czego służy opisuje obrazy, kontenery i porty. Do prawdziwego stacku uruchom brokera razem z aplikacją przez Compose — zobacz Czym Jest Docker Compose i Jak Go Używać:
| |
worker dociera do brokera pod hostname’em rabbitmq. Compose stawia oba serwisy w tej samej sieci user-defined, więc nazwa serwisu sama działa jako rozwiązywalny hostname — mechanizm opisany w Co to Jest Sieć Docker i Jak Jej Używać.
Publikowanie i odbieranie wiadomości
pika to standardowy klient Pythona dla protokołu RabbitMQ (AMQP 0-9-1). Zainstaluj go przez pip install pika. Producer publikujący zadanie wygląda tak:
| |
durable=True w queue_declare mówi RabbitMQ, żeby zachować samą kolejkę po restarcie brokera. delivery_mode=pika.DeliveryMode.Persistent mówi, żeby zapisywać każdą wiadomość na dysk, zamiast trzymać ją tylko w pamięci. Bez obu tych ustawień docker restart na brokerze po cichu kasuje wszystko, co jeszcze nie zostało odebrane.
Consumer pobiera wiadomości z tej samej kolejki i potwierdza każdą dopiero, gdy praca jest naprawdę skończona:
| |
Dwa ustawienia robią tu całą robotę związaną z niezawodnością:
basic_qos(prefetch_count=1)blokuje RabbitMQ przed wysłaniem workerowi drugiej wiadomości, zanim potwierdzi pierwszą. Bez tego zajęty worker może siedzieć na dziesięciu wiadomościach, podczas gdy wolny nie dostaje nic.basic_ackpo pracy, nie przed nią. Jeśli proces padnie w połowie zmniejszania obrazu, wiadomość nigdy nie została potwierdzona, więc RabbitMQ dostarcza ją ponownie do innego workera zamiast ją stracić. Zostawauto_ackwyłączone, co i tak jest domyślnym ustawieniem pika.
Uruchom trzy kopie skryptu consumera, a RabbitMQ rozdzieli kolejkę między nimi metodą round-robin. To cała historia skalowania w tym wzorcu: żadnego kodu koordynującego, żadnego współdzielonego stanu między workerami.
Niektóre wiadomości zawodzą za każdym razem, gdy są dostarczane: źle sformułowany payload, wywołanie do serwisu, którego już na stałe nie ma. Ponowne dostarczanie zamienia to w nieskończoną pętlę przez twoich consumerów. Ustaw x-dead-letter-exchange na kolejce, a wiadomości, które odrzucasz bez ponownego kolejkowania (basic_nack z requeue=False), trafiają do osobnego exchange, gdzie sprawdzasz je ręcznie.
Jak exchange’y RabbitMQ kierują wiadomości
Przykłady powyżej publikują z exchange="", czyli default exchange, który kieruje wiadomość prosto do kolejki podanej w routing_key. To pokrywa większość przypadków kolejek zadań. Prawdziwy model RabbitMQ stawia exchange przed każdą kolejką, a zmiana typu exchange zmienia sposób routingu:
| Typ exchange | Kieruje do | Użyj do |
|---|---|---|
direct (default exchange jest jednym z nich) | Kolejki dokładnie pasującej do routing key | Kolejek zadań, jedna grupa consumerów na kolejkę |
fanout | Każdej podpiętej kolejki, ignorując routing key | Rozgłaszania jednego zdarzenia do kilku niezależnych consumerów |
topic | Kolejek, których wzorzec bindowania pasuje do routing key (order.*.created) | Selektywnego rozgłaszania, gdy consumerzy chcą podzbiór typów zdarzeń |
Exchange typu fanout to sposób na uzyskanie publish/subscribe w RabbitMQ: publikujesz order_placed raz, a serwis email, serwis analityczny i serwis wykrywania oszustw dostają każdy swoją kopię przez własną kolejkę, zamiast rywalizować o te same wiadomości.
Kiedy kolejka wiadomości to zły wybór
- Wywołujący potrzebuje odpowiedzi, żeby sam odpowiedzieć. Kolejka służy do “zrób to kiedyś”, nie do “policz to teraz”. Jeśli twój endpoint checkoutu potrzebuje policzonego kosztu wysyłki, zanim odpowie, to jest wywołanie synchroniczne, nie zadanie w kolejce.
- Potrzebujesz ścisłej kolejności w całej kolejce. RabbitMQ gwarantuje kolejność na poziomie kolejki przy jednym consumerze, ale dodaj drugiego dla przepustowości, a dwie wiadomości mogą skończyć w złej kolejności. Jeśli kolejność ma znaczenie, jak w logu zdarzeń czy maszynie stanów, potrzebujesz narzędzia zbudowanego pod uporządkowane strumienie, jak partycje Kafki kluczowane po ID encji.
- Zadanie jest małe, synchroniczne, i nie masz jeszcze brokera. Wątek w tle albo job runner in-process to mniejsza powierzchnia operacyjna niż stawianie i monitorowanie RabbitMQ dla zadania wielkości crona.
- Masz już Redisa i tolerujesz sporadyczną utratę zadania. Redis Streams albo biblioteka jak BullMQ dają ci kolejkę bez nowej infrastruktury, kosztem słabszych gwarancji dostarczenia przy awarii brokera. Sprawdź najpierw, co Redis już robi w twoim stacku: Jak Działa Cache Redis i Jak Go Używać.
Przenieś pierwsze zadanie do kolejki
Wybierz coś, co blokuje request bez potrzeby: email z potwierdzeniem, miniaturkę, ping do analityki. Uruchom RabbitMQ w Dockerze, zadeklaruj dla tego zadania jedną trwałą kolejkę, przenieś je. Ustaw prefetch_count=1 i potwierdzaj ręcznie od pierwszego dnia, bo dodawanie niezawodności po tym, jak awaria już pożarła jedną wiadomość, kosztuje o wiele więcej niż napisanie tych dwóch linijek teraz. Sięgnij po exchange fanout albo topic dopiero, gdy masz drugiego consumera, który chce tego samego zdarzenia.