Пайплайн веб-данных настолько надёжен, насколько надёжен его уровень сбора. Когда дашборды устаревают или цифры перестают сходиться, причина почти никогда не в аналитическом коде. Проблема в начале пайплайна: скрапер, сломавшийся после редизайна сайта, запросы, которые начали блокироваться, или страницы, которые открываются в браузере, но возвращают пустую оболочку на обычный HTTP-запрос. Если воспринимать сбор данных как нечто хрупкое, весь пайплайн наследует эту хрупкость.

В этом руководстве показано, как создать масштабируемый пайплайн веб-данных, где Crawlbase отвечает за сбор, а стандартные ETL-инструменты за остальное. Вы будете собирать страницы с помощью Crawling API для работы по требованию и асинхронного Crawler для высокопроизводительных задач, трансформировать и валидировать «сырой» HTML, загружать чистые строки в хранилище и планировать всё это с мониторингом. Каждый шаг содержит работающий код, который можно адаптировать.

Как выглядит масштабируемый пайплайн веб-данных

Паттерн представляет собой классическую форму ETL с одним важным разграничением зон ответственности. Crawlbase занимает позицию в начале как уровень приёма данных и берёт на себя всё, что делает скрапинг нестабильным: рендеринг JavaScript, ротацию IP, маршрутизацию запросов и защиту от блокировок. Ваши системы занимаются парсингом, валидацией, хранением и аналитикой. Поток данных идёт слева направо:

bash
Web  ->  Crawlbase (collect)  ->  Transform + Validate  ->  Storage  ->  BI / ML

Смысл проведения границы именно здесь, в долговечности. Внешние сайты не являются стабильными зависимостями; они выпускают изменения разметки, проводят эксперименты и развёртывают защиту от ботов без предупреждения. Размещение управляемого уровня сбора данных в начале превращает изменение сайта в вопрос конфигурации, а не в аварию пайплайна. Crawlbase предоставляет два инструмента сбора для двух типов нагрузки, и production-пайплайн обычно использует оба.

  • Crawling API для получения данных в реальном времени по требованию для известных URL. Вы отправляете URL, он возвращает страницу.
  • Async Crawler для крупномасштабного сбора в режиме «отправил и забыл». Вы добавляете URL, он загружает их асинхронно и отправляет результаты на ваш вебхук через POST.

Это то же разграничение, к которому приходит любая серьёзная операция скрапинга для ecommerce: быстрый путь для точечных запросов и пакетный путь для широкого охвата. Если вы только знакомитесь с механикой прокси, статья что такое прокси-сервер даст полезный контекст, хотя смысл управляемого API в том, что вам не нужно ничем управлять самостоятельно.

Шаг 1: Сбор данных через Crawling API

Crawling API принимает URL и ваш токен, а затем возвращает отрендеренную страницу. Вы отправляете HTTP GET; он маршрутизирует запрос через пул ротирующих IP, опционально рендерит JavaScript при использовании JS-токена и возвращает HTML (или разобранный JSON). Простейший вызов, одиночный curl:

bash
curl 'https://api.crawlbase.com/?token=YOUR_TOKEN&url=https%3A%2F%2Fexample.com%2Fproducts'

В пайплайне нужен небольшой переиспользуемый коллектор вместо «сырого» curl. Установите официальный клиент и оберните вызов так, чтобы остальная часть пайплайна получала чистый HTML и никогда не думала о токенах или повторных попытках. Используйте JS-токен для страниц с клиентским рендерингом и обычный токен для статического HTML.

bash
python3 -m venv venv && source venv/bin/activate
pip install crawlbase
python
from crawlbase import CrawlingAPI

api = CrawlingAPI({'token': 'YOUR_TOKEN'})

def collect(url, render=False):
    options = {'ajax_wait': True, 'page_wait': 3000} if render else {}
    response = api.get(url, options)
    status = response['status_code']
    if status != 200:
        raise RuntimeError(f'collect failed for {url}: {status}')
    return response['body'].decode('utf-8')

html = collect('https://example.com/products', render=True)
print(len(html), 'bytes')

Два момента делают этот код уровня production, а не прототипом. Во-первых, он проверяет status_code и выдаёт исключение при всём, что не является чистой загрузкой, поэтому плохая страница выдаёт явную ошибку вместо того, чтобы отравлять хранилище пустыми строками. Во-вторых, флаг render делает вызовы честными в отношении того, какие страницы требуют JavaScript: платите за рендеринг только там, где контент действительно его требует. Этот коллектор, та единица, которую ваш планировщик будет вызывать для каждого известного URL.

Normal token vs JS token

Crawlbase выдаёт вам два токена. Обычный токен быстро и дёшево возвращает статический HTML; JS-токен сначала рендерит страницу в реальном браузере, что необходимо для сайтов с клиентским рендерингом. Используйте JS-токен только когда страница возвращает пустую оболочку на обычный запрос, и сочетайте его с ajax_wait и page_wait, чтобы поздно загружаемый контент успел появиться.

Шаг 2: Масштабирование объёма с помощью async Crawler

Crawling API синхронный: один запрос, один ответ, и ваш код ждёт. Это именно то, что нужно для нескольких сотен известных URL. Для десятков тысяч блокировка на каждом вызове не масштабируется. Асинхронный Crawler меняет модель. Вы добавляете URL в именованный краулер, запрос немедленно возвращает Request ID, Crawlbase загружает страницу в фоне, и по завершении отправляет результат на ваш callback-эндпоинт через POST. Ничего в вашем коде не блокируется в ожидании страниц.

Вы переходите в асинхронный режим, добавив два параметра к тому же эндпоинту: callback=true и crawler=YourCrawlerName (краулер создаётся один раз в панели управления и направляется на ваш URL вебхука). Добавление URL выглядит так:

bash
curl 'https://api.crawlbase.com/?token=YOUR_TOKEN&callback=true&crawler=my-pipeline&url=https%3A%2F%2Fexample.com%2Fp%2F123'

Вместо тела страницы вы получаете Request ID, что означает постановку URL в очередь:

json
{ "rid": "1e92e8bf4618772871c14d4" }

С вашей стороны отправка большого пакета, это плотный цикл. Суть в пропускной способности: вы запускаете все URL, не ожидая завершения ни одного из них, а очередь поглощает нагрузку.

python
from crawlbase import CrawlingAPI

api = CrawlingAPI({'token': 'YOUR_TOKEN'})

def push_batch(urls):
    options = {'callback': True, 'crawler': 'my-pipeline'}
    for url in urls:
        response = api.get(url, options)
        rid = response['body']['rid']
        print(f'queued {url} -> {rid}')

push_batch([
    'https://example.com/p/123',
    'https://example.com/p/124',
    'https://example.com/p/125',
])

Вторая часть асинхронного режима, обработчик callback. Crawlbase отправляет скрапнутую страницу на вебхук, зарегистрированный вами в краулере, помещая HTML в тело запроса и метаданные (Request ID, исходный URL и статус) в заголовки. Ваш обработчик должен делать минимум: быстро ответить кодом 200 и передать полезную нагрузку на шаг трансформации. Тяжёлый парсинг в обработчике рискует вызвать таймаут доставки и повторные попытки.

javascript
const express = require('express')
const app = express()

// Crawlbase POSTs raw HTML; capture the body as text
app.use(express.text({ type: '*/*', limit: '10mb' }))

app.post('/crawlbase/callback', (req, res) => {
  const rid = req.headers['rid']
  const url = req.headers['url']
  const status = req.headers['original_status']

  // ack immediately, process out of band
  res.sendStatus(200)

  enqueueForTransform({ rid, url, status, html: req.body })
})

app.listen(8080, () => console.log('callback listening on :8080'))

Если не хотите запускать вебхук вообще, направьте краулер на Crawlbase Cloud Storage и опрашивайте его вместо этого; компромисс, небольшая задержка в обмен на отсутствие инфраструктуры. В любом случае асинхронная модель позволяет собирать миллионы страниц, не блокируя ваше приложение ни на одной загрузке.

Crawlbase Crawling API + Crawler

Один токен покрывает обе половины сбора: синхронные вызовы для известных URL и асинхронные запросы для больших объёмов, с рендерингом, ротацией IP и защитой от блокировок на стороне сервера. Начните с бесплатного тарифа, подключите callback к тестовому эндпоинту и наблюдайте, как результаты поступают ещё до того, как вы построите остальную часть пайплайна.

Шаг 3: Трансформация и валидация «сырого» HTML

Сбор данных даёт вам HTML. Шаг трансформации превращает этот HTML в чистые типизированные записи и отбрасывает всё, что не проходит проверку качества. Именно здесь многие пайплайны незаметно деградируют: задание сообщает об успехе, но записанные строки пустые, потому что селектор устарел. Выполняйте валидацию явно, чтобы сбой парсинга выглядел как сбой.

Парсите любым подходящим инструментом; в примере используется BeautifulSoup. Функция извлекает поля, нормализует их в нативные типы и отказывается создавать запись с отсутствующим названием или непарсируемой ценой.

python
import re
from bs4 import BeautifulSoup

def transform(html, source_url):
    soup = BeautifulSoup(html, 'html.parser')
    records = []

    for card in soup.select('.product-card'):
        name = card.select_one('.title')
        price = card.select_one('.price')
        if not name or not price:
            continue  # skip incomplete cards, do not emit junk

        digits = re.sub(r'[^\d.]', '', price.get_text())
        if not digits:
            continue

        records.append({
            'name': name.get_text(strip=True),
            'price': float(digits),
            'source_url': source_url,
        })

    if not records:
        raise ValueError(f'no records parsed from {source_url} (selectors may have drifted)')

    return records

Ключевая структура: преобразуйте каждое поле в нативный тип (цена как float, строка без пробелов), пропускайте неполные записи вместо записи пустых значений и генерируйте исключение, когда целая страница не даёт никаких результатов, чтобы устаревший селектор обнаруживался в тот же день, когда он сломался, а не через недели в отчёте. Если хотите полностью пропустить парсинг для поддерживаемых сайтов, Crawling API возвращает структурированный JSON напрямую, и этот шаг становится транзитным.

Шаг 4: Загрузка в хранилище

Имея валидированные записи, запишите их туда, где можно делать запросы. Место назначения зависит от масштаба и использования: реляционная база данных PostgreSQL для транзакционного доступа, хранилище вроде BigQuery для аналитики, поисковое хранилище или потоковая платформа для нижестоящих потребителей. SQLite достаточно для демонстрации паттерна, и этот паттерн и является тем, что обобщается: выполняйте upsert по стабильному ключу, чтобы повторный запуск пайплайна обновлял существующие строки вместо их дублирования.

python
import sqlite3
from datetime import datetime, timezone

def load(records, db_path='pipeline.db'):
    conn = sqlite3.connect(db_path)
    conn.execute('''
        CREATE TABLE IF NOT EXISTS products (
            source_url TEXT PRIMARY KEY,
            name TEXT NOT NULL,
            price REAL NOT NULL,
            collected_at TEXT NOT NULL
        )''')

    now = datetime.now(timezone.utc).isoformat()
    for r in records:
        conn.execute('''
            INSERT INTO products (source_url, name, price, collected_at)
            VALUES (?, ?, ?, ?)
            ON CONFLICT(source_url) DO UPDATE SET
                name=excluded.name,
                price=excluded.price,
                collected_at=excluded.collected_at
        ''', (r['source_url'], r['name'], r['price'], now))

    conn.commit()
    conn.close()

Upsert делает шаг загрузки идемпотентным: повторный запуск одного пакета оставляет таблицу в том же состоянии, что именно то, что нужно, когда планировщик повторяет неудавшийся запуск. Метка времени collected_at даёт сигнал актуальности, который используется для мониторинга на следующем шаге. Замените вызовы SQLite на клиент вашего хранилища, и логика перенесётся без изменений.

Шаг 5: Автоматизация, планирование и мониторинг

Компоненты складываются в одну функцию пайплайна, которую и вызывает планировщик. Связывание сбора, трансформации и загрузки с блоком try/except для каждого URL не позволяет одной плохой странице убить весь запуск.

python
import logging

logging.basicConfig(level=logging.INFO)
log = logging.getLogger('pipeline')

def run_pipeline(urls):
    ok, failed = 0, 0
    for url in urls:
        try:
            html = collect(url, render=True)
            records = transform(html, url)
            load(records)
            ok += 1
        except Exception as err:
            failed += 1
            log.error('pipeline failed for %s: %s', url, err)

    log.info('run complete: %d ok, %d failed', ok, failed)
    if failed > ok:
        raise RuntimeError('majority of URLs failed, check upstream')

Для запуска по расписанию проще всего использовать cron. Эта запись запускает пайплайн каждые шесть часов и добавляет вывод в лог, который можно отслеживать или отправлять в систему мониторинга:

bash
# run the pipeline every 6 hours
0 */6 * * * /path/to/venv/bin/python /path/to/run.py >> /var/log/pipeline.log 2>&1

Cron отлично подходит для небольшого числа задач. Когда появляются зависимости между шагами, повторные попытки и ретроспективные заполнения, переходите к оркестратору рабочих процессов, Apache Airflow или Prefect, которые предоставляют DAG, автоматические повторы и интерфейс для истории запусков. Для async Crawler на стороне сбора данных вообще не нужен планировщик: вы добавляете URL, а результаты поступают в callback по мере их готовности.

Мониторинг, это разница между пайплайном, которому вы доверяете, и тем, за которым постоянно следите. Отслеживайте минимум три вещи. Объём: количество строк за каждый запуск, чтобы внезапное снижение сигнализировало о проблеме со сбором. Актуальность: метки времени collected_at, которые вы сохранили, чтобы можно было отправлять оповещение, когда данные устаревают. Частота сбоев: соотношение успешных и неудачных запросов из каждого запуска, чтобы постепенное увеличение предупреждало об изменении целевого сайта до того, как всё сломается. Дополните это разумными практиками скрапинга; статья как скрапить сайты без блокировки охватывает практики, поддерживающие уровень сбора данных в здоровом состоянии при масштабировании.

Итоги

Ключевые выводы

  • Сбор данных, слабое звено. Поместите управляемый уровень приёма в начало, чтобы изменение сайта было вопросом конфигурации, а не аварией пайплайна.
  • Два режима сбора. Crawling API обслуживает синхронные запросы для известных URL; async Crawler обрабатывает большие объёмы и отправляет результаты на ваш вебхук без блокировки.
  • Валидация в трансформации. Преобразуйте поля в нативные типы, пропускайте неполные записи и генерируйте исключение, когда страница не даёт результатов, чтобы устаревшие селекторы давали явный сбой.
  • Сделайте загрузку идемпотентной. Выполняйте upsert по стабильному ключу, чтобы повторы и повторные запуски обновляли строки вместо их дублирования.
  • Планируйте и мониторьте. Cron или оркестратор управляет запусками; отслеживайте объём, актуальность и частоту сбоев, чтобы выявлять проблемы заранее.

Часто задаваемые вопросы

Как создать масштабируемый пайплайн веб-данных с Crawlbase?

Используйте Crawlbase как уровень сбора и стандартные ETL-инструменты для остального. Собирайте страницы через Crawling API для известных URL и через async Crawler для больших объёмов, трансформируйте возвращаемый HTML в валидированные типизированные записи, загружайте их в хранилище с идемпотентным upsert и планируйте запуски через cron или оркестратор, отслеживая объём, актуальность и частоту сбоев. Crawlbase берёт на себя рендеринг, ротацию IP и защиту от блокировок, поэтому ваш код работает только с чистыми данными.

Когда использовать Crawling API, а когда async Crawler?

Используйте Crawling API, когда у вас есть известный список URL и нужна страница немедленно, это подходит для бэкенд-сервисов, задач мониторинга и запросов в реальном времени. Используйте async Crawler при сборе больших объёмов или когда нужна доставка в режиме «отправил и забыл»: вы добавляете URL, мгновенно получаете Request ID, а Crawlbase отправляет каждый результат на ваш callback по завершении. Многие пайплайны используют оба: API для точечного получения и Crawler для широкого охвата.

Как работает callback async Crawler?

Вы создаёте именованный краулер в панели управления и направляете его на URL вашего вебхука, затем добавляете URL с параметрами callback=true и crawler=YourCrawlerName. Каждое добавление немедленно возвращает Request ID. Когда Crawlbase заканчивает загрузку страницы, он отправляет HTTP POST на ваш вебхук с HTML в теле и метаданными в заголовках. Ваш обработчик должен быстро вернуть 200 и обрабатывать полезную нагрузку асинхронно, чтобы доставка не истекла по таймауту.

Нужно ли самостоятельно управлять прокси или защитой от ботов?

Нет. Crawling API и Crawler направляют запросы через ротирующий пул IP, рендерят JavaScript при использовании JS-токена и применяют защиту от блокировок на стороне сервера. Вы отправляете URL и получаете страницу, минуя необходимость держать собственный пул прокси и флот headless-браузеров. Если вам нужны только «сырые» ротирующие IP для вашего стека, Smart AI Proxy предоставляет ту же сеть как стандартный прокси-эндпоинт.

Как не допустить записи пустых или некорректных данных?

Выполняйте валидацию на шаге трансформации. Проверяйте статус ответа при сборе и генерируйте исключение при всём, что не является чистой загрузкой; затем при парсинге пропускайте записи с отсутствующими обязательными полями и генерируйте исключение, когда целая страница даёт ноль записей, поскольку это обычно означает устаревший селектор. Делайте загрузку идемпотентной через upsert, чтобы повторные попытки не дублировали строки, и сохраняйте метку времени сбора для мониторинга актуальности и оповещений при устаревании данных.

Может ли этот пайплайн обрабатывать миллионы страниц?

Да. Узким местом в наивной реализации является блокировка на каждом синхронном запросе, которую async Crawler устраняет, ставя работу в очередь и доставляя результаты через callback. Добавляйте большие пакеты не ожидая, пусть очередь поглощает нагрузку и обрабатывайте результаты по мере поступления. Для очень крупных или непрерывных программ план Enterprise добавляет пропускную способность и поддержку, необходимые для высокопроизводительного сбора данных.

Начать создавать

Обходите любой сайт в масштабе, без борьбы с инфраструктурой.

Crawlbase берёт на себя прокси, отпечатки и CAPTCHA, чтобы ваша команда выпускала конвейеры данных вместо поддержки обвязки краулинга. 1 000 запросов бесплатно, без карты.

Самообслуживание · Звонок отдела продаж не требуется · Доступны корпоративные объёмы краулинга