Ми не раз бачили, як парсинг у циклі розвалюється при першій же помилці мережі. Пропадають тисячі рядків, а налагодження займає години. Черга завдань вирішує три фундаментальні проблеми: ізоляція збоїв, автоматичні повторні спроби та горизонтальне масштабування воркерів. В одному з проектів з 50 000 сторінок каталогу ми перейшли з лінійного скрипта на BullMQ — час виконання скоротився вдвічі, а відсоток втрачених даних впав до нуля. Середня затримка відновлення після збою становить 60 секунд завдяки експоненціальним ретраям.
Яку проблему вирішує черга?
Замість послідовного обходу ви ставите завдання в чергу і забуваєте. Якщо воркер впав, завдання повертається в чергу і повторюється з експоненціальною затримкою. Потоки паралельної обробки налаштовуються через 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% від бюджету на парсинг.







