Строим бронепоезд очередей

Строим бронепоезд очередей

Когда вы строите архитектуру проекта, где предполагается много сообщений между разными компонентами, вы начинаете хотеть брокеры очередей. Так размышлял и я для своего OpenSource проекта Gracehub. Однако, решил пойти своим путём.

Что максимально важно: по возможности "переложить" надежность доставки сообщений с сети на сам сервер - создать очередь и гарантию(!) доставки. Мне так же важна низкая сложность старта моего проекта как по ресурсам VPS так и навыкам деплойщика.

А теперь внимание вопрос: что дешевле, оперативка или ssd? Особенно когда проект не высоконагружен, а на долгосроке vps кушает ваши деньги. Именно такими умозаключениями я пришел к тому, что мне не подходят:
1. n8n - не нужен затратный ui да еще + Redis как брокер(уже два сервиса надо!), я хочу чтобы запуск моего проекта был доступен на самом слабом железе.
2. RabbitMQ - Заяц Кроликович хорош, но это дополнительный внешний сервис на хосте с кворумом + publisher confirms это сопоставимо с постгри на небольшой нагрузке очередей. К тому же у меня нет сложных правил маршрутизации.

Поэтому было принято решение написать свой масштабируемый(1 воркер очереди = 1 контейнер) микромодуль очередей для хранения в postgre.
Плюсы:
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. База сама "будит" воркеров:

# queue_worker.py - настройка подписки listen_conn = await db.pool.acquire() await listen_conn.add_listener('tg_update_channel', on_notify) # Триггер в БД при вставке новой задачи: CREATE OR REPLACE FUNCTION notify_new_tg_update() RETURNS TRIGGER AS $ BEGIN PERFORM pg_notify('tg_update_channel', NEW.instance_id); # <-= Сигнал воркерам! RETURN NEW; END; $ LANGUAGE plpgsql;

Вместо 1000 запросов/сек на polling (при 10 воркерах с таймаутом 0.01с) — 0 запросов в idle-режиме! Работать не от забора до обеда а умно!

4. Идемпотентный retry-механизм - надо обезопасить зону зависших и сбойных задач. Обработка неудач:

# database.py - обработка неудач async def fail_tg_update(self, job_id: int, error: str, *, max_attempts: int = 10): await self.execute(""" UPDATE tg_update_queue SET status = CASE WHEN attempts < $2::int THEN 'retry' ELSE 'dead' END, run_at = CASE WHEN attempts < $2::int THEN NOW() + ($3::int * interval '1 second') ELSE run_at END, last_error = $4::text, locked_at = NULL, # <-=- анлок для повторной попытки locked_by = NULL WHERE id = $1::bigint RETURNING status """, (job_id, max_attempts, retry_seconds, error))

Задачи могут безопасно возвращаться в очередь:

# queue_worker.py - автоматическое восстановление cleanup_service = QueueCleanupService(db) cleanup_service.start() #<-- Запуск фонового мониторинга # Внутри сервиса периодически выполняется: UPDATE tg_update_queue SET status = 'retry', run_at = NOW(), locked_at = NULL, last_error = COALESCE(last_error, '') || 'stuck requeued' WHERE status = 'processing' AND locked_at IS NOT NULL AND locked_at < NOW() - ($1 * INTERVAL '1 second')

Даже если воркер упал во время обработки, задача вернётся в очередь.

ACID-гарантии "из коробки"

Самое мощное преимущество - "бронированная" транзакционность за счет того, что PostgreSQL умеет из коробки. Мы обеспечиваем гарантию доставки сообщения, как только оно "коснулось" нашего не раздутого по ресурсам VPS.

❗Важно: UNLOGGED не делаю для таблицы очередей. Хотя это ускорило бы бд, но ломает репликацию - очередь не доедет до реплики(Нет записей в WAL).

Когда НЕ использую эту архитектуру:

  1. >10k сообщений/сек - уже нужен будет уже брокер. Однако, такого бомболейло сервера надо ещё добиться большим кол-вом инстансов моих пользователей.
  2. Сложная маршрутизация - тут её нет, но n8n при сложных workflow очень полюбился коллегам по цеху. КроликMQ для приверженцев старой школы.
  3. Очереди >100 млн сообщений - нужно партицировать(пьём воду глотками). В текущей схеме это пока не предусмотрено.

Посмотрим, как покажут себя полевые нагрузки.

Ссылки на проект: Github GitVerse Канал