Мы не раз видели, как парсинг в цикле разваливается при первой же ошибке сети. Пропадают тысячи строк, а отладка занимает часы. Очередь задач решает три фундаментальных проблемы: изоляция сбоев, автоматические повторные попытки и горизонтальное масштабирование воркеров. В одном из проектов с 50 000 страниц каталога мы перешли с линейного скрипта на BullMQ — время выполнения сократилось вдвое, а процент потерянных данных упал до нуля. Средняя задержка восстановления после сбоя составляет 60 секунд благодаря экспоненциальным ретраям. Бюджет на реализацию очереди задач обычно составляет от $300–1k в зависимости от сложности и объёма данных.
Какую проблему решает очередь?
Вместо последовательного обхода вы ставите задачу в очередь и забываете. Если воркер упал, задача возвращается в очередь и повторяется по экспоненциальной задержке. Потоки параллельной обработки настраиваются через concurrency, а при росте нагрузки добавляются новые инстансы. Типичная конфигурация для среднего проекта — 5–10 воркеров с concurrency 5, что даёт до 50 одновременных задач.
Как выбрать брокер очереди?
Для большинства веб-проектов оптимальны BullMQ или Celery. BullMQ работает на Redis и предоставляет UI Board для мониторинга, поддерживает приоритеты и до 100 тысяч задач в сутки на одном инстансе. Celery лучше подходит для Python-стека: цепочки задач и групповые обработки строятся без дополнительного кода. RabbitMQ оправдан в высоконагруженных системах, где требуется сложная маршрутизация через routing keys и гарантированная доставка на уровне AMQP. Например, при агрегации данных из 20+ источников с разной скоростью. Официальная документация RabbitMQ рекомендует DLQ для критичных данных.
Сравните характеристики:
| Характеристика | BullMQ | Celery | RabbitMQ |
|---|---|---|---|
| Бэкенд | Redis | Redis/RabbitMQ | AMQP |
| Макс. throughput | ~100k/сутки | ~50k/сутки | >200k/сутки (кластер) |
| Встроенный UI | Да (Board) | Flower | Да (Management) |
| Сложность настройки | Низкая | Средняя | Высокая |
BullMQ: настройка воркера и ретраи
import { Queue, Worker, Job } from 'bullmq'; import { Redis } from 'ioredis'; const connection = new Redis({ host: 'localhost', port: 6379, maxRetriesPerRequest: null }); // Создание очереди export const scrapeQueue = new Queue('scraping', { connection, defaultJobOptions: { attempts: 3, backoff: { type: 'exponential', delay: 60_000 }, removeOnComplete: { count: 500 }, removeOnFail: { count: 200 }, }, }); // Добавление задачи await scrapeQueue.add('scrape-url', { url: 'https://example.com/catalog?page=5', siteId: 42, depth: 1, }, { priority: 1 }); // Воркер const worker = new Worker('scraping', async (job: Job) => { const { url, siteId } = job.data; const html = await fetchWithProxy(url); const products = parseProducts(html); await saveProducts(products, siteId); return { count: products.length }; }, { connection, concurrency: 5 }); worker.on('failed', (job, err) => { logger.error(`Job ${job?.id} failed: ${err.message}`); }); Celery: пайплайн с цепочками задач
from celery import Celery, chain, chord import redis app = Celery('scraper', broker='redis://localhost:6379/0', backend='redis://localhost:6379/1') app.conf.task_routes = { 'scraper.tasks.fetch_listing': {'queue': 'listings'}, 'scraper.tasks.fetch_product': {'queue': 'products'}, } @app.task(bind=True, max_retries=3, default_retry_delay=60) def fetch_listing(self, url: str, site_id: int) -> list[str]: try: html = fetch_page(url) return extract_product_urls(html) except (NetworkError, RateLimitError) as exc: raise self.retry(exc=exc, countdown=2 ** self.request.retries * 60) @app.task(bind=True, max_retries=3) def fetch_product(self, url: str, site_id: int) -> dict: try: html = fetch_page(url) return parse_product(html) except Exception as exc: raise self.retry(exc=exc) @app.task def save_products(products: list[dict], site_id: int): bulk_upsert(products, site_id) # Запуск пайплайна def start_site_crawl(site_id: int, catalog_url: str): urls = fetch_listing.delay(catalog_url, site_id).get() chord( fetch_product.s(url, site_id) for url in urls )(save_products.s(site_id)) Пример конфигурации Celery с rate limiting
app.conf.task_annotations = { 'scraper.tasks.fetch_product': { 'rate_limit': '10/m' } } Это ограничивает количество задач fetch_product до 10 в минуту на воркер, что помогает избежать блокировки по IP.
Dead Letter Queue: настройка и анализ
Задачи, исчерпавшие все попытки, попадают в Dead Letter Queue. Это не просто мусорная корзина — это очередь для ручного анализа и переработки. В RabbitMQ DLQ настраивается через аргументы очереди:
channel.queue_declare( queue='scraping.products', durable=True, arguments={ 'x-dead-letter-exchange': 'scraping.dlx', 'x-dead-letter-routing-key': 'failed', 'x-message-ttl': 3600000, # 1 час } ) channel.exchange_declare(exchange='scraping.dlx', exchange_type='direct') channel.queue_declare(queue='scraping.failed', durable=True) channel.queue_bind(queue='scraping.failed', exchange='scraping.dlx', routing_key='failed') Задачи из DLQ можно переотправить в основную очередь после устранения причины сбоя — через Admin UI или скрипт. В BullMQ DLQ реализуется с помощью отдельной очереди и обработчика failed.
Сравнение стратегий повторных попыток
| Стратегия | Задержка | Когда использовать |
|---|---|---|
| Экспоненциальная | 2^retry * base | При временных ошибках сети |
| Линейная | retry * base | При rate limiting |
| Постоянная | fixed delay | При стабильных условиях |
Мониторинг очереди
BullMQ Board (UI для BullMQ) или Flower (для Celery) дают визуальное представление о состоянии очередей. Ключевые метрики, которые стоит отслеживать:
- Глубина очереди (waiting jobs)
- Скорость обработки (jobs/sec)
- Процент ошибок по типам задач
- Время выполнения (p50, p95, p99)
Эти метрики экспортируются в Prometheus через /metrics эндпоинт и визуализируются в Grafana. Среднее время ответа наших парсеров после внедрения мониторинга снизилось на 35%.
Процесс работы
- Аналитика: определяем объём данных, частоту парсинга, требования к надёжности.
- Проектирование: выбираем брокер, схему задач, настройки ретраев.
- Реализация: пишем воркеры, DLQ, интеграция с мониторингом.
- Тестирование: прогоняем нагрузку, проверяем поведение при сбоях.
- Деплой: разворачиваем на сервере или в kubernetes, подключаем CI/CD.
Что входит в работу
- Настройка очереди (BullMQ/Celery/RabbitMQ) с политикой ретраев и DLQ
- Интеграция с Redis (Sentinel/Cluster) или RabbitMQ
- Мониторинг (Prometheus + Grafana) и алерты
- Документация по эксплуатации и обучение команды
- Гарантия работоспособности в течение месяца после сдачи
Мы занимаемся парсинговыми системами 6+ лет, реализовали более 30 проектов с очередями. Получите консультацию инженера — мы оценим ваш проект за один день и предложим оптимальную архитектуру.
Сроки реализации
Базовая очередь с ретраями и DLQ — 3–4 рабочих дня. Добавление метрик, UI и кластеризации — ещё 2–3 дня. Итоговая стоимость рассчитывается индивидуально под ваш объём данных.
Свяжитесь с нами — мы предложим решение под вашу задачу. Экономия на повторных ошибках может достигать 40% от бюджета на парсинг.







