Модуль 20

Streams у RabbitMQ (огляд)

Ось твій урок про RabbitMQ Streams, написаний у стилі лекції CS50. Вмикай уяву, ми починаємо!


🎓 Тема: RabbitMQ Streams (Огляд)

"Це не просто черга, це — історія"


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

Привіт усім! Радий вас бачити. Давайте почнемо з простого питання.

Уявіть, що ви надсилаєте повідомлення другу в месенджері, але там діє правило: як тільки друг прочитав повідомлення — воно зникає назавжди. Немає історії переписки, немає можливості прокрутити вгору і згадати, про що ви домовлялися вчора.

Саме так працюють класичні черги в RabbitMQ (Classic/Quorum Queues). Логіка проста: 1. Producer (Видаввець) кладе задачу в чергу. 2. Consumer (Споживач) забирає її. 3. Pouf! 💥 Задача зникає з пам'яті брокера.

Це чудово працює для розподілу завдань (наприклад, обробка замовлення піци). Ви зробили піцу, віддали кур'єру — вам не треба тримати її на кухні.

Але виникає проблема. 🛑 А що, коли до вашої системи приєднується новий сервіс (наприклад, "Аналітика"), і він хоче проаналізувати всі замовлення за останній місяць? У класичній черзі це неможливо — повідомлення вже "з'їли" інші сервіси. Вони зникли.

Риторичне питання: Чи не було б чудово мати механізм, який працює як стрічка новин у Facebook або лог-файл? Де повідомлення лежать, ми можемо їх читати, перечитувати, і навіть "відмотувати час назад"?

Ось тут на сцену виходять RabbitMQ Streams.


2. 🧠 Теоретична база (Як це працює під капотом)

Забудьте на хвилину слово "Черга" (Queue). Давайте думати про Журнал (Log).

RabbitMQ Streams моделюють структуру даних, яку називають "Append-only Log" (Журнал тільки для допису).

⚙️ Ключові концепти (Без нудних визначень):

  1. Append-only (Тільки додавання): Нові повідомлення завжди пишуться в кінець "стрічки". Ми ніколи не змінюємо і не видаляємо повідомлення з середини.
  2. Immutability (Незмінність): Як тільки повідомлення записане — воно висічене в камені.
  3. Offset (Зміщення/Вказівник):
    • У звичайній черзі споживач "з'їдає" повідомлення.
    • У Stream споживач просто рухає свій "палець" (offset) по стрічці. "Я прочитав повідомлення №5, тепер читаю №6".
    • Важливо: Якщо Споживач А прочитав повідомлення №5, воно не зникає. Споживач Б може прийти через годину і теж прочитати повідомлення №5.
  4. Retention Policy (Політика зберігання): Оскільки ми нічого не видаляємо при читанні, диск може переповнитися. Тому ми кажемо RabbitMQ: "Зберігай історію за останні 7 днів" або "Зберігай останні 50 ГБ даних". Старі дані автоматично відрізаються з "хвоста".

Аналогія з життя: 📼 * Classic Queue: Це прямий ефір телебачення. Пропустив момент — він пропав. * Stream: Це Netflix або запис на відеокасету. Ти можеш дивитися новинки, а можеш натиснути "rewind" і подивитися серію, яка вийшла тиждень тому.


3. 🧪 Приклади (Код і Логіка)

Давайте подивимось, як це виглядає на практиці. Для прикладів використаємо псевдокод, максимально наближений до Python/Java, щоб зрозуміти суть.

Приклад 1: Публікація (Producer)

Тут все майже як завжди. Ми просто пишемо в Stream.

stream_name = "user-clicks-stream"

# Ми просто додаємо події в кінець журналу
publisher.send(stream_name, "User A clicked Button X")
publisher.send(stream_name, "User B clicked Button Y")
publisher.send(stream_name, "User A clicked Logout")

Очікування: Повідомлення лежать у RabbitMQ під номерами 0, 1, 2.

Приклад 2: Читання в реальному часі (Consumer Live)

Цей споживач хоче бачити, що відбувається прямо зараз.

consumer = Consumer(stream_name)
consumer.set_offset("next") # Читати тільки нове

consumer.start()
# Чекає нових повідомлень...

Приклад 3: Подорож у часі (Replay) 🕰️

А ось це — "кілер-фіча". Уявіть, що у вас впала база даних, і вам треба відновити стан системи, перечитавши всі події з самого початку.

backup_service = Consumer(stream_name)

# МАГІЯ ТУТ 👇
backup_service.set_offset("first") # Почати з найпершого доступного повідомлення!

for msg in backup_service:
    print(f"Відновлюю подію: {msg.body}")
    # Процес пройдеться по всім повідомленням: 0, 1, 2... і наздожене реальний час.

Чому це круто? Тому що RabbitMQ Stream оптимізований для неймовірно швидкого читання з диска. Він використовує sendfile (системний виклик ОС), щоб перекидати байти з диска в мережу без зайвого навантаження на процесор. Це дуже швидко! 🚀


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

Ви маєте закріпити це. Уявіть, що ви архітектор системи для великого інтернет-магазину.

Завдання 1: "Налаштування" Уявіть, що ви створюєте Stream в інтерфейсі RabbitMQ. * Ви ставите max-age (час життя) = "24h". * Що станеться з повідомленням, яке прийшло 25 годин тому? * (Відповідь: воно буде видалено системою автоматично, навіть якщо його ніхто не прочитав).

Завдання 2: "Новий співробітник" У вас є Stream orders, де зберігаються всі замовлення. Працює сервіс Shipping (Відправка), який читає події в реальному часі (offset: last). Раптом маркетологи просять підключити сервіс BonusPoints, щоб нарахувати бонуси за всі покупки за останній тиждень. * Як ви налаштуєте offset для сервісу BonusPoints? * Чи завадить це роботі сервісу Shipping?

Завдання 3: "Помилка в коді" Ви написали код, який читає Stream, обробляє повідомлення і... падає з помилкою на повідомленні №100500. Ви перезапускаєте сервіс. * Якщо ви не зберігали свій прогрес (offset) десь у базі, звідки почне читати сервіс? * (Підказка: за замовчуванням він може почати знову з початку або з кінця, залежно від налаштувань. У Streams важливо "комітити" або зберігати свій offset, щоб знати, де ви зупинилися).

Завдання 4: Міні-кейс Придумайте сценарій для IoT (Інтернет речей). У вас є 1000 датчиків температури, які шлють дані щосекунди. * Чому тут краще використати Stream, а не Classic Queue? (Подумайте про швидкість запису і про те, що дані потрібні кільком споживачам: для моніторингу в реальному часі і для архіву).


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

Як досвідчений інженер, я хочу застерегти вас від типових помилок.

1. Не плутайте Streams з Kafka. Так, RabbitMQ Streams натхненні Kafka. Але RabbitMQ простіший у піднятті та адмініструванні для багатьох задач. Проте, якщо у вас петабайти даних — можливо, Kafka все ще краще. Думайте: "Мені потрібен Smart Proxy (Rabbit) чи Dumb Pipe (Kafka)?". Streams роблять RabbitMQ чимось посереднім.

2. Диск — не гумовий. Новачки часто забувають налаштувати Retention Policy. * Помилка: Створити стрім без лімітів. * Результат: Через місяць сервер падає, бо диск забитий на 100%. * Порада профі: Завжди ставте ліміт або по часу (наприклад, 7 днів), або по розміру (наприклад, 50GB).

3. Stream — це не про "Task Queue". Якщо вам треба, щоб повідомлення було видалено відразу після обробки, і щоб його точно не обробили двічі випадково — використовуйте старі добрі Classic Queues. Streams — це про дані, а не про задачі.


6. 🧩 Підсумок

Отже, що ми сьогодні дізналися?

  1. RabbitMQ Streams — це журнал (log), який тільки дописується.
  2. Дані не зникають після читання. Вони зникають, коли спливає їхній термін (retention).
  3. Ми можемо підключати нових споживачів, які читатимуть старі дані, не заважаючи іншим.
  4. Це ідеально для високого навантаження (High Throughput) і патерну "Fan-out" (коли одні дані потрібні багатьом сервісам).

Тепер ви вмієте: Розрізняти, коли використовувати "Трубу" (Stream), а коли "Поштову скриньку" (Queue). Ви можете спроектувати систему, яка не боїться втратити історію подій.

🔜 У наступному уроці: Ми поговоримо про Sharding та Super Streams. Що робити, коли потік води (даних) настільки потужний, що одна труба (стрім) просто трісне? Ми будемо вчитися розділяти потоки!

А поки що — це був CS50... тобто RabbitMQ Streams! 😉