Как устроена очередь сообщений и зачем она нужна, если есть база
Рано или поздно в любом проекте всплывает одна и та же боль: пользователь нажал кнопку, а сервер вместо мгновенного ответа завис на пять секунд, пока где-то в фоне отправляется письмо или перекодируется видео. Первая мысль — «а давайте просто запишем задачу в таблицу базы данных и будем её периодически проверять». Мысль рабочая, но у неё есть предел прочности, и в этой статье разберём, где он проходит и почему очередь сообщений — не модная прослойка ради красивой архитектуры, а конкретное решение конкретных проблем.
Содержание
- Что такое очередь сообщений и из каких частей она состоит
- Что происходит с сообщением от отправки до обработки
- Почему просто «пишите в таблицу и читайте её» — не то же самое
- RabbitMQ и Kafka — два разных подхода к одной идее
- Практический сценарий: тяжёлая задача не должна блокировать ответ пользователю
- Быстрый пример: очередь на RabbitMQ в docker compose
- Куда очередь ведёт дальше: мониторинг, ретраи, dead-letter
Что такое очередь сообщений и из каких частей она состоит
В основе — три роли и одна структура данных:
- Продюсер (producer) — часть системы, которая создаёт сообщение и отправляет его в очередь. Это может быть веб-бэкенд, который принял запрос пользователя, cron-задача, вебхук от внешнего сервиса.
- Брокер (broker) — отдельный сервис (RabbitMQ, Kafka, NATS, Redis Streams), который принимает сообщение, сохраняет его и решает, кому и когда его отдать.
- Консьюмер (consumer) — процесс, который забирает сообщение из очереди и обрабатывает его: отправляет письмо, режет видео на превью, дёргает внешний API.
- Очередь (queue) или топик (topic) — упорядоченная структура внутри брокера, куда складываются сообщения и откуда их забирают консьюмеры.
Продюсер и консьюмер никогда не общаются друг с другом напрямую — они оба разговаривают только с брокером. Это ключевое отличие от синхронного HTTP-вызова: продюсер положил сообщение и тут же освободился, ему не нужно ждать, пока консьюмер закончит работу. Он даже не обязан знать, сколько консьюмеров сейчас живо и живы ли они вообще.
Именно это называется развязкой по времени и по скорости. Продюсер может генерировать тысячу сообщений в секунду, а консьюмер — обрабатывать по одному в секунду; очередь просто накапливает буфер, и обе стороны продолжают работать в своём темпе, не блокируя друг друга.
Что происходит с сообщением от отправки до обработки
Механика внутри брокера примерно одинакова у большинства систем, различаются детали:
- Продюсер отправляет сообщение брокеру — обычно это TCP-соединение с определённым протоколом (AMQP у RabbitMQ, свой бинарный протокол у Kafka).
- Брокер записывает сообщение на диск (или в память, в зависимости от настройки durability) и подтверждает продюсеру приём — это называется publisher confirm.
- Сообщение попадает в очередь и ждёт, пока свободный консьюмер его заберёт. В RabbitMQ это обычно push-модель: брокер сам отправляет сообщение подписанному консьюмеру. В Kafka — pull-модель: консьюмер сам опрашивает партицию и читает следующее смещение (offset).
- Консьюмер обрабатывает сообщение и отправляет брокеру подтверждение — ack. Только после ack брокер считает сообщение доставленным и может его удалить (RabbitMQ) или сдвинуть офсет консьюмера (Kafka).
- Если консьюмер упал или не ответил за отведённое время, брокер возвращает сообщение обратно в очередь (RabbitMQ) либо просто не двигает офсет (Kafka) — сообщение будет обработано ещё раз при следующем подключении.
Последний пункт — источник важного нюанса: большинство брокеров по умолчанию гарантируют at-least-once delivery — сообщение доставится минимум один раз, но иногда доставится дважды (например, консьюмер обработал задачу, но упал до отправки ack). Отсюда практическое правило: обработчик сообщений должен быть идемпотентным — повторная обработка того же сообщения не должна приводить к дублированию результата (второе письмо клиенту, повторное списание). Обычно это решается уникальным ключом операции, который проверяется перед выполнением действия.
Нужен сервер под эту задачу?
Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.
Арендовать серверПочему просто «пишите в таблицу и читайте её» — не то же самое
Идея хранить очередь задач прямо в базе выглядит соблазнительно: одна инфраструктура, никакого нового сервиса, знакомый SQL. Проблема не в том, что так нельзя — при небольшой нагрузке это рабочий паттерн (его называют transactional outbox или просто job table). Проблема в том, какие сложности вылезают, когда нагрузка и число консьюмеров растут.
Нет нативного уведомления о новой записи. База данных — это пассивное хранилище: она не умеет сама сказать «эй, появилась новая строка, обработай её». Консьюмеру приходится опрашивать таблицу по таймеру:
SELECT * FROM jobs WHERE status = 'pending' ORDER BY created_at LIMIT 10;
Это называется polling, и у него есть цена. Если опрашивать раз в секунду — задача может ждать до секунды даже когда очередь пустая, а база получает лишнюю нагрузку постоянными запросами вхолостую. Если опрашивать реже — растёт задержка обработки. Брокер сообщений вместо этого либо сам пушит сообщение консьюмеру, либо консьюмер блокируется на чтении сокета и просыпается ровно в момент появления данных — без пустых циклов.
Конкурентный доступ нескольких консьюмеров сложнее гарантировать без блокировок. Если у вас два воркера читают одну и ту же таблицу задач, оба легко могут выбрать одну и ту же строку и обработать её дважды. Решение существует и в PostgreSQL — конструкция SELECT ... FOR UPDATE SKIP LOCKED:
SELECT id, payload FROM jobs
WHERE status = 'pending'
ORDER BY created_at
LIMIT 1
FOR UPDATE SKIP LOCKED;
Это рабочий приём (на нём построены некоторые лёгкие очереди на Postgres), но он требует, чтобы разработчик сам продумал транзакции, таймауты блокировок, обработку зависших задач при падении воркера посреди транзакции. Брокер сообщений эту логику уже реализовал и проверил на миллионах инсталляций — вам не нужно отлаживать свою версию distributed lock.
Очередь оптимизирована именно под паттерн доставки, а не под произвольные запросы. База данных хороша в том, для чего она проектировалась — гибкие выборки, джойны, транзакционная целостность связанных данных. Она не оптимизирована под миллион вставок и удалений в секунду в одну таблицу — это упирается в разрастание индексов, autovacuum, конкуренцию за одну и ту же горячую страницу. Брокер сообщений, наоборот, спроектирован ровно под один паттерн: быстро принять, быстро отдать, гарантировать порядок (FIFO) внутри очереди или партиции.
Итог этого раздела не «таблица — плохо, очередь — хорошо», а честная граница: для одного фонового воркера и сотен задач в час таблица в Postgres будет проще в эксплуатации, чем отдельный брокер. Как только типов задач и независимо масштабируемых консьюмеров становится несколько — специализированный брокер отбивает свою сложность.
RabbitMQ и Kafka — два разных подхода к одной идее
Оба решают задачу «продюсер — брокер — консьюмер», но с разной философией, и путать их не стоит:
| RabbitMQ | Kafka | |
|---|---|---|
| Модель | Традиционная очередь: сообщение удаляется после ack | Журнал (log): сообщение хранится retention-период независимо от чтения |
| Доставка | Push брокера консьюмеру | Pull консьюмера у партиции по офсету |
| Повторное чтение | Нет — сообщение исчезло после обработки | Да — можно перечитать историю, сбросив офсет |
| Типичный сценарий | Задачи (отправить письмо, обработать заказ) | Потоки событий, аналитика, множество независимых читателей одних данных |
| Порядок | Гарантирован в рамках очереди | Гарантирован в рамках партиции |
Если задача — «выполнить действие один раз и забыть» (отправить email, сгенерировать PDF), классическая очередь RabbitMQ ближе к сути задачи. Если нужно, чтобы несколько разных систем могли независимо читать один и тот же поток событий (аналитика, аудит, репликация в другую систему) — там сильнее Kafka с её моделью журнала. Подробнее сравнение с NATS — ещё одной лёгкой альтернативой — разбирали в статье NATS или Kafka: что выгоднее и когда.
Практический сценарий: тяжёлая задача не должна блокировать ответ пользователю
Классический пример, ради которого всё это затевается. Пользователь регистрируется на сайте и должен получить email с подтверждением. Наивная реализация:
POST /register
→ создать пользователя в базе
→ отправить email через SMTP (может занять 1-3 секунды, а может и таймаутнуть)
→ вернуть 200 OK
Если SMTP-сервер провайдера подвиснет или ответит с задержкой, пользователь будет смотреть на крутящийся индикатор загрузки все эти секунды — а при таймауте получит ошибку, хотя аккаунт уже создан. С очередью сценарий меняется:
POST /register
→ создать пользователя в базе
→ положить сообщение {user_id, email} в очередь "send-welcome-email"
→ вернуть 200 OK (заняло миллисекунды)
отдельный процесс-consumer:
→ читает очередь "send-welcome-email"
→ отправляет email, при ошибке — retry с задержкой
→ если после N попыток не вышло — сообщение уходит в dead-letter queue для ручного разбора
Тот же принцип работает для обработки видео (загрузили файл — ответили пользователю мгновенно, а перекодирование в разные разрешения идёт в фоне и статус потом подтягивается через вебсокет или поллинг статуса), генерации отчётов, массовой рассылки, вызовов внешних API с непредсказуемой задержкой. Общий признак задач-кандидатов на очередь: результат не нужен пользователю прямо в ответе на его запрос, и/или операция может занять заметное время или иногда падать.
Отдельный практический плюс: consumer и продюсер (веб-сервер) масштабируются независимо. Если писем стало в 10 раз больше — добавляете ещё воркеров-консьюмеров, не трогая веб-слой. Если веб-трафик вырос, а писем не прибавилось — масштабируете только веб-часть.
Быстрый пример: очередь на RabbitMQ в docker compose
Чтобы не оставаться на уровне теории — минимальный рабочий стенд. Файл docker-compose.yml:
version: "3.8"
services:
rabbitmq:
image: rabbitmq:3-management
ports:
- "5672:5672" # AMQP-порт для приложений
- "15672:15672" # веб-панель управления
environment:
RABBITMQ_DEFAULT_USER: app
RABBITMQ_DEFAULT_PASS: changeme
volumes:
- rabbitmq_data:/var/lib/rabbitmq
volumes:
rabbitmq_data:
Продюсер на Python (библиотека pika):
import pika
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost', credentials=pika.PlainCredentials('app', 'changeme'))
)
channel = connection.channel()
channel.queue_declare(queue='send-welcome-email', durable=True)
channel.basic_publish(
exchange='',
routing_key='send-welcome-email',
body='{"user_id": 42, "email": "user@example.com"}',
properties=pika.BasicProperties(delivery_mode=2) # сохранить на диск
)
connection.close()
Консьюмер:
import pika
def callback(ch, method, properties, body):
print(f"Обрабатываю: {body}")
# тут реальная отправка письма
ch.basic_ack(delivery_tag=method.delivery_tag) # подтверждаем только после успеха
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost', credentials=pika.PlainCredentials('app', 'changeme'))
)
channel = connection.channel()
channel.queue_declare(queue='send-welcome-email', durable=True)
channel.basic_qos(prefetch_count=1) # не забирать следующее, пока не подтвердили текущее
channel.basic_consume(queue='send-welcome-email', on_message_callback=callback)
channel.start_consuming()
Обратите внимание на basic_qos(prefetch_count=1) — без этого RabbitMQ по умолчанию может раздать консьюмеру сразу пачку сообщений, и если процесс упадёт, все неподтверждённые вернутся в очередь разом. И на delivery_mode=2 — без него сообщения хранятся только в памяти и пропадут при перезапуске брокера.
Панель управления доступна на http://<ip-сервера>:15672 — там видно глубину очереди и скорость обработки, что полезно на старте, чтобы своими глазами увидеть, накапливается очередь или консьюмер успевает её разгребать. Развернуть Kafka на VPS тем же способом разобрано в статье как установить и настроить Kafka на VPS.
Куда очередь ведёт дальше: мониторинг, ретраи, dead-letter
Три вещи, о которых стоит подумать сразу, а не когда очередь уже легла в проде:
- Глубина очереди как метрика. Если число сообщений в очереди стабильно растёт — консьюмеры не успевают за продюсером: либо добавляйте воркеров, либо ищите, почему обработка стала медленнее.
- Ретраи с задержкой, а не мгновенный повтор. Если внешний сервис (SMTP, платёжный шлюз) временно недоступен, повторять запрос немедленно бессмысленно — лучше выдержать растущую паузу между попытками (экспоненциальный backoff), иначе вы сами создадите нагрузку на тот же сервис своими ретраями.
- Dead-letter queue (DLQ). Отдельная очередь для сообщений, которые не удалось обработать после N попыток. Это не мусорка, а место, куда человек приходит вручную посмотреть, что пошло не так, вместо того чтобы терять данные молча.
Если фоновые задачи в проекте пока живут в отдельной таблице через cron — переход к очереди не обязан быть резким. Настройку самих периодических задач разбирали в статье как установить и настроить cron-задачи на VPS; часто первый шаг — вынести из cron-скрипта именно "толстую" операцию в очередь, оставив cron лишь диспетчером.
Нужен сервер под эту задачу?
Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.
Арендовать серверНужны сами нейросети для контента?
Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.
Частые вопросы
Очередь сообщений — это то же самое, что кеш вроде Redis?
Нет. Redis умеет выступать простым брокером (списки, Pub/Sub, Redis Streams), но его основное назначение — хранилище ключ-значение в памяти. Если у вас уже есть Redis в проекте и нагрузка невысокая, Redis Streams — разумный старт; для сложных гарантий доставки и больших объёмов обычно берут специализированный брокер. Про установку Redis отдельно — в статье как установить и настроить Redis на VPS.
Можно ли гарантировать, что сообщение обработается ровно один раз (exactly-once)?
На уровне транспорта — практически никогда полностью, честнее говорить об at-least-once плюс идемпотентность на стороне консьюмера. Некоторые системы (Kafka с транзакциями между чтением и записью в саму Kafka) приближаются к exactly-once внутри своей экосистемы, но как только в цепочке появляется внешняя система (SMTP, платёжный шлюз), гарантия снова упирается в идемпотентность приложения.
Что будет, если брокер упадёт?
Зависит от настроек durability. Если очереди и сообщения объявлены durable (как в примере выше) и брокер пишет на диск, после перезапуска накопленные неподтверждённые сообщения останутся и будут обработаны. Если брокер работал только в памяти без persistence — при падении содержимое очереди теряется, поэтому для важных данных durability отключать нельзя.
Сколько ресурсов сервера нужно под очередь?
Зависит от объёма сообщений и настроек retention (особенно у Kafka, где хранится журнал, а не только неподтверждённые сообщения). Для старта с невысокой нагрузкой обычно достаточно скромной конфигурации в 1-2 vCPU и 2 ГБ RAM — но это ориентир, а не гарантия: реальная потребность считается по размеру сообщений, их числу в секунду и периоду хранения, который вы задаёте.
Нужно ли сразу брать Kafka, если проект небольшой?
Как правило нет. Kafka создавалась под большие потоки событий и множество независимых консьюмеров одних и тех же данных; она тяжелее в эксплуатации, чем RabbitMQ. Для типового веб-приложения с фоновыми задачами (письма, отчёты, вебхуки) RabbitMQ или даже очередь на Postgres через SKIP LOCKED закрывают задачу с меньшими эксплуатационными расходами.
Обсудить статью, задать вопрос или начать новую тему
Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.
Перейти в сообщество →