Строим бронепоезд очередей
Когда вы строите архитектуру проекта, где предполагается много сообщений между разными компонентами, вы начинаете хотеть брокеры очередей. Так размышлял и я для своего OpenSource проекта Gracehub. Однако, решил пойти своим путём.
Что максимально важно: по возможности "переложить" надежность доставки сообщений с сети на сам сервер - создать очередь и гарантию(!) доставки. Мне так же важна низкая сложность старта моего проекта как по ресурсам VPS так и навыкам деплойщика.
А теперь внимание вопрос: что дешевле, оперативка или ssd? Особенно когда проект не высоконагружен, а на долгосроке vps кушает ваши деньги. Именно такими умозаключениями я пришел к тому, что мне не подходят:
1. n8n - не нужен затратный ui да еще + Redis как брокер(уже два сервиса надо!), я хочу чтобы запуск моего проекта был доступен на самом слабом железе. 2. RabbitMQ - Заяц Кроликович хорош, но это дополнительный внешний сервис на хосте с кворумом + publisher confirms это сопоставимо с постгри на небольшой нагрузке очередей. К тому же у меня нет сложных правил маршрутизации.
Поэтому было принято решение написать свой масштабируемый(1 воркер очереди = 1 контейнер) микромодуль очередей для хранения в постгри.
Плюсы: 1. Меньше расходы на ram вашего vps. Успешный запуск на слабом железе 2. Максимальная надежность, благодаря ACID транзакциям
Проблемы, которые нужно было решить: 1. Блокируемость 2. IO нагрузку 3. Скорость работы
Как я решал это: 1. Сделал SKIP LOCKED для записей этой таблицы - воркеры должны переходить к следующей записи, зная что предыдущая уже обрабатывается другим воркером: ``` #метод pick_tg_update( WITH cte AS ( SELECT id FROM tg_update_queue WHERE status IN ('pending', 'retry') AND run_at <= NOW() ORDER BY run_at ASC, id ASC FOR UPDATE SKIP LOCKED # <-- Не мешать другим воркерам! LIMIT 1 ) UPDATE tg_update_queue q SET status = 'processing', attempts = q.attempts + 1, locked_at = NOW(), locked_by = $1 FROM cte WHERE q.id = cte.id_ RETURNING q.*;_ ``` Преимущество: N воркеров могут одновременно выбирать N разных задач без взаимных блокировок.
2. Частичный индекс - индекс на только необработанные сообщения в очередь. Чтобы быстро отсортировать тех, с кем надо работать Индекс строил только для необработанных задач, что резко сокращает его размер и ускоряет выборку: ``` CREATE INDEX IF NOT EXISTS idx_tg_update_queue_pending_active ON tg_update_queue (run_at, id) WHERE status IN ('pending', 'retry') # <-- Только активные задачи! ``` Если у вас 1 млн выполненных задач и 1000 активных, индекс в 1000 раз меньше полного!
3. LISTEN/NOTIFY вместо Polling SELECT запросов в бд. Превращаем базу в активного брокера очередей и экономим железо vps. База сама "будит" воркеров.
4. Идемпотентный retry-механизм - надо обезопасить зону зависших и сбойных задач. Обработка неудач.
ACID-гарантии "из коробки Самое мощное преимущество - "бронированная" транзакционность за счет того, что PostgreSQL умеет из коробки. Мы обеспечиваем гарантию доставки сообщения, как только оно "коснулось" нашего не раздутого по ресурсам VPS.
❗Важно: UNLOGGED не делаю для таблицы очередей. Хотя это ускорило бы бд, но ломает репликацию - очередь не доедет до реплики(Нет записей в WAL).
Когда НЕ использую эту архитектуру >10k сообщений/сек - уже нужен будет уже брокер. Однако, такого бомболейло сервера надо ещё добиться большим кол-вом инстансов моих пользователей.
Сложная маршрутизация - тут её нет, но n8n при сложных workflow очень полюбился коллегам по цеху. КроликMQ для приверженцев старой школы.
Очереди >100 млн сообщений - нужно партицировать(пьём воду глотками). В текущей схеме это пока не предусмотрено.
Посмотрим, как покажут себя полевые нагрузки. Ещё многое нужно доделать.
GitHub - https://github.com/glebgv/GraceHub GitVerse - https://gitverse.ru/glebgv/GraceHub