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







