Модуль 21

RabbitMQ у Python (pika)

Ось твій урок у стилі CS50! 🔥


🎓 Тема уроку: RabbitMQ у Python (бібліотека pika)

Привіт, друзі! 👋

Сьогодні ми розберемо технологію, яка перетворює "повільні" та "неповороткі" програми на швидкі, масштабовані та надійні системи. Ми говоримо про Message Brokers (брокери повідомлень), а саме — про RabbitMQ.


1. 🔥 Вступ: Проблема «завислого вікна»

Уявіть, що ви зайшли на сайт інтернет-магазину, замовили новий ноутбук і натиснули кнопку «Купити». Крутиться значок завантаження... ⏳ Проходить 1 секунда... 2 секунди... 5 секунд... 10 секунд...

Ви починаєте нервувати. "Що сталося? Гроші списали? Інтернет зник?". А насправді в цей момент сервер магазину намагається відправити вам PDF-чек на email, оновити склад, повідомити службу доставки та записати статистику. І поки він це все не зробить — він не скаже вам "Дякую за замовлення".

Питання до вас: Чи повинен клієнт чекати, поки сервер відправить email? Звісно, ні! Клієнт хоче побачити "Все ок!" миттєво.

Як це вирішити? Аналогія з рестораном 🍽️

Уявіть гарний ресторан. * Ви (Клієнт) робите замовлення. * Офіціант (Web-сервер/Producer) записує його.

Якби офіціант працював як поганий код, він би взяв замовлення, пішов на кухню, став біля плити і чекав, поки кухар приготує страву, щоб одразу винести її вам. Поки він чекає — він не може обслуговувати інших.

У реальності офіціант робить інакше: 1. Бере замовлення. 2. Кріпить папірець на спеціальну рейку (Черга/Queue) на кухні. 3. Миттєво повертається до зали приймати нові замовлення. 4. Кухар (Worker/Consumer) бачить папірець, готує страву, коли звільниться.

RabbitMQ — це і є та сама "рейка" для папірців на стероїдах. Це спосіб сказати: "Зроби це, але я не буду чекати над душею".


2. 🧠 Теоретична база (Що під капотом?)

RabbitMQ — це Message Broker (посередник повідомлень). Він приймає повідомлення, зберігає їх і віддає тому, хто має їх обробити.

У Python ми працюємо з ним через бібліотеку pika.

🔑 Ключові поняття (Це треба запам'ятати):

  1. Producer (Продюсер/Виробник) — програма, яка надсилає повідомлення (ваш веб-сайт).
  2. Queue (Черга) — поштова скринька всередині RabbitMQ, де лежать повідомлення. Вона живе в пам'яті або на диску.
  3. Consumer (Споживач/Воркер) — програма, яка чекає на повідомлення і обробляє їх (скрипт, що шле пошту).

Як це працює логічно:

Потік даних виглядає так: Producer ---> [ RabbitMQ (Queue) ] ---> Consumer

Інтуїтивно: Це як конвеєр. Ви кидаєте задачу на стрічку конвеєра і йдете далі. Десь там в кінці стоїть робітник, який цю задачу підбере і виконає.

Чому не просто записати в базу даних?

Бо RabbitMQ вміє: * Балансувати навантаження (роздавати задачі 10-м воркерам). * Гарантувати доставку (якщо воркер впав — задача не зникне, а піде іншому). * З'єднувати різні мови (сайт на Python, воркер на Java).


3. 🧪 Приклади: Від Hello World до Реальності

Для роботи вам знадобиться встановлений RabbitMQ (найпростіше через Docker) та pip install pika.

Приклад 1: "Hello World" (Надсилання)

Уявімо, що це наш Офіціант. Він просто каже "Привіт".

Файл: send.py (Producer)

import pika

# 1. Встановлюємо з'єднання з RabbitMQ (він має бути запущений локально)
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# 2. Створюємо чергу 'hello'.
# Це важливо: якщо черги немає, повідомлення полетить в нікуди!
channel.queue_declare(queue='hello')

# 3. Відправляємо повідомлення
channel.basic_publish(exchange='',
                      routing_key='hello', # Ім'я черги
                      body='Привіт, RabbitMQ!')

print(" [x] Sent 'Привіт, RabbitMQ!'")

# 4. Закриваємо з'єднання
connection.close()

Приклад 2: Отримання (Consumer)

А це наш Кухар. Він сидить і чекає.

Файл: receive.py (Consumer)

import pika, sys, os

def main():
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()

    # Знову оголошуємо чергу. Навіщо?
    # Бо ми не знаємо, хто запуститься першим: send.py чи receive.py.
    # Це гарантує, що черга існує.
    channel.queue_declare(queue='hello')

    # Функція, яка спрацює, коли прийде повідомлення (callback)
    def callback(ch, method, properties, body):
        print(f" [x] Отримано: {body.decode()}")

    # Підписуємось на чергу
    channel.basic_consume(queue='hello',
                          auto_ack=True, # Поки що автоматично підтверджуємо отримання
                          on_message_callback=callback)

    print(' [*] Чекаю повідомлень. Натисніть CTRL+C для виходу')
    channel.start_consuming() # Нескінченний цикл

if __name__ == '__main__':
    try:
        main()
    except KeyboardInterrupt:
        print('Interrupted')
        try:
            sys.exit(0)
        except SystemExit:
            os._exit(0)

Що ви очікуєте побачити? Запустіть receive.py в одному терміналі. Він зависне в очікуванні. Відкрийте інший термінал і запустіть send.py. Бах! У першому вікні миттєво з'явиться текст. Магія! ✨


Приклад 3: Реальна задача (Work Queues)

А тепер ускладнимо. Уявіть, що обробка кожного повідомлення займає час (наприклад, ресайз картинки). Ми хочемо, щоб повідомлення не губилися, якщо воркер "впаде" під час роботи.

Тут вмикається механізм Message Acknowledgment (підтвердження).

Змінений worker.py:

import pika
import time

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

channel.queue_declare(queue='task_queue', durable=True) # durable=True - щоб черга пережила рестарт RabbitMQ

def callback(ch, method, properties, body):
    print(f" [x] Отримано задачу: {body.decode()}")

    # Імітуємо важку роботу (кількість крапок = секунди роботи)
    time.sleep(body.count(b'.')) 

    print(" [x] Виконано")

    # ! ВАЖЛИВО: Ручне підтвердження
    # Кажемо RabbitMQ: "Я закінчив, можеш видаляти повідомлення з черги"
    ch.basic_ack(delivery_tag=method.delivery_tag)

# channel.basic_qos(prefetch_count=1) # "Не давай мені нову задачу, поки я не зроблю поточну"
channel.basic_consume(queue='task_queue', on_message_callback=callback)

print(' [*] Чекаю задач...')
channel.start_consuming()

4. 🛠 Практична частина

Час забруднити руки кодом! Виконайте ці завдання:

  1. 🔹 Запуск: Запустіть приклад 1 і 2. Переконайтеся, що вони "спілкуються".
  2. 🔹 Експеримент з JSON: Змініть send.py так, щоб він відправляв не просто рядок, а JSON-об'єкт (наприклад, {"user_id": 1, "email": "test@test.com"}). У receive.py розпарсіть цей JSON і виведіть лише email.
    • Підказка: Використовуйте модуль json.
  3. 🔹 Краш-тест (Hardcore):
    • Візьміть код із Прикладу 3 (без auto_ack=True, з ручним підтвердженням).
    • Відправте задачу, яка "робиться" 10 секунд.
    • Поки воркер пише, що працює — вбийте його (Ctrl+C).
    • Запустіть воркер знову.
    • Питання: Чи прийшла задача знову? (Має прийти, бо ви не відправили ack!).
  4. 🔹 Міні-кейс: Напишіть систему "Спам-фільтр".
    • producer.py відправляє коментарі користувачів.
    • consumer.py перевіряє, чи є в тексті слово "bad". Якщо є — пише в консоль "BAN!", якщо немає — "OK".

5. 💡 Мислення як у розробника

Як досвідчений інженер дивиться на черги?

  1. "А що, якщо воркер помре?" Новачки часто забувають про ACK (acknowledgment). Якщо ви використовуєте auto_ack=True і ваш скрипт впаде з помилкою посеред обробки — повідомлення зникне назавжди. Досвідчені завжди використовують ручний ack в кінці успішної обробки (try...finally).

  2. Проблема "Отруйного повідомлення" (Poison Pill). Уявіть повідомлення, яке завжди викликає помилку у вашому коді. Воркер бере його -> падає -> RabbitMQ бачить, що воркер впав -> повертає повідомлення в чергу -> воркер встає -> бере його знову -> падає. Це нескінченний цикл! Порада: Робіть обробку помилок і відкладайте "біті" повідомлення в окрему "чергу мерців" (Dead Letter Queue).

  3. Ідемпотентність (Страшне слово, проста суть). Що, якщо RabbitMQ подумав, що воркер впав, і віддав задачу іншому, а насправді перший воркер вже списав гроші з картки, але просто затримався з відповіддю? Ви спишете гроші двічі. Порада: Пишіть код так, щоб повторна обробка тієї ж задачі не ламала систему (наприклад, перевіряйте статус замовлення в БД перед списанням).


6. 🧩 Підсумок

Сьогодні ми з вами: * ✅ Зрозуміли, навіщо потрібні черги (щоб не змушувати користувача чекати). * ✅ Навчилися відправляти повідомлення через pika. * ✅ Навчилися приймати та обробляти їх. * ✅ Дізналися про важливість ack для надійності.

Тепер ви можете будувати системи, де один шматок коду не блокує інший. Ваші веб-сайти будуть літати! 🚀

Тизер: Але що, як ми хочемо, щоб одне повідомлення отримали одразу всі воркери (наприклад, оновити кеш на всіх серверах)? Або відфільтрувати повідомлення за темами? Про це ми поговоримо на наступному уроці: Exchanges (Fanout, Direct, Topic).

А поки — кодіть і не блокуйте головний потік! 😉