大規模ウェブスクレイピングアーキテクチャ: フィールドガイド

数百ページならスクリプト。数百万ページなら分散システム。ここでは、キュー、非同期ワーカー、プロキシ階層化、リトライ、重複排除、データ品質チェックを含むアーキテクチャを紹介します。実際に動くコード付きです。

数百ページのスクレイピングはスクリプトです。数百万ページのスクレイピングは分散システムです。ターゲット数が「ノートパソコンで一晩で実行」から「今週中に完了させる必要がある」に変わると、難しい部分は解析ではなく、それを取り巻くすべてになります。作業をキューに入れ、分散させ、ブロックを避け、失敗をリトライし、結果を保存してチェックする方法です。これはアーキテクチャとしての大規模ウェブスクレイピングであり、スニペットではなく、各段階で実際に動くコードがあります。

スケールとは同時実行性であり、より大きなループではありません

一つの数字が具体化します。2万のリスティングページがあり、各ページに20アイテムがあるカテゴリーを考えてみましょう。40万ページを取得する必要があります。現実的な2.5秒/ページで、厳密な順次実行は約1,000,000秒、つまり11.5日のページロード待ち時間です。200ページを並行して処理すると、その11.5日は1時間の実時間に縮まります。スケールでの制約は時間であり、同時実行性がそれを取り戻す方法です。アーキテクチャの他のすべては、その同時実行性を生き残らせるために存在します。

400,000ページが順次で11.5日かかるのに対し、200を並行して行うと約1時間で済むという統計図。100万件中1%の失敗で10,000件の失敗が発生。
同時実行性は11.5日の実行を1時間に変えますが、100万リクエストでは1%の失敗率でも10,000ページが失敗します。

アーキテクチャの概要

数百万ページを生き残るスクレイパーは、小規模な分散システムであり、ボリュームでのみ現れる問題を解決するいくつかの名前付きパーツを持っています。

まずキュー: 発見と取得を分離する

最も重要な構造的決定は、「何をスクレイプするか」と「スクレイプを実行すること」の間にキューを置くことです。プロデューサーがURLを列挙し、ワーカーのプールがそれを消化します。どちらの側も他がどれだけ速く動くかを知りませんし、プロデューサーに手を加えずにワーカーを追加できます。PythonではこれはRedis上のCeleryまたはRQです。NodeではBullMQ、大規模ではRabbitMQまたはKafkaです。1ファイルでのパターン:

import asyncio, aiohttp

CONCURRENCY = 50
queue = asyncio.Queue()

async def worker(session):
    while True:
        url = await queue.get()
        try:
            async with session.get(url, timeout=20) as resp:
                await handle(url, await resp.text(), resp.status)
        except Exception as err:
            await on_failure(url, err)
        finally:
            queue.task_done()

async def run(urls):
    for u in urls:
        queue.put_nowait(u)
    async with aiohttp.ClientSession() as session:
        tasks = [asyncio.create_task(worker(session)) for _ in range(CONCURRENCY)]
        await queue.join()
        for t in tasks:
            t.cancel()

重要なつまみはCONCURRENCYです。低すぎるとスケールを可能にする並行性が無駄になります。高すぎるとターゲットと自分の出口を圧倒します。エラー率が上昇するのを見て適切な値を見つけます。それがまさに監視がシステムの一級の部分である理由であり、後から考えることではありません。

大規模スクレイピングパイプラインのフローダイアグラム: キュー、非同期ワーカー、プロキシ層、解析と品質チェック、ストレージ
8つの名前付きパーツ。プロキシ、アンチボット、レンダリングレイヤーは健康を保つのが最も難しく、購入する自然な場所です。

プロキシ層が最初に壊れる

低ボリュームではアンチボット防御にほとんど気づきませんが、スケールでは最初に実行を壊します。1つのIPから数十万のリクエストを送信すると、レート制限され、その後チャレンジされ、最終的にブロックされます。修正は多くのアドレスにわたる回転です。階層化します: 寛容なターゲットとAPIには安価なデータセンターIPを、実際のユーザートラフィックを期待する商業ターゲットには回転する住宅IPを使用します。しかし、回転だけでは不十分です。現代の防御はTLSフィンガープリントとヘッダー順序も読み取るため、トラフィックは新しいIPから来るだけでなく、ブラウザのように見えなければなりません。

リトライ: 失敗は通常状態

100万リクエストでは1%の一時的な失敗率でも10,000ページが失敗します。このボリュームでは失敗はエッジケースではなく、パイプラインは失敗した取得を致命的ではなく通常のものとして扱わなければなりません。指数バックオフとキャップを使ってリトライし、URLをデッドレターキューに移動して実行をブロックしないようにします。なぜ失敗したのかを読み取ります: タイムアウトや503はリトライする価値がありますが、ハード404はそうではありません。

import asyncio, random

async def fetch_with_retry(session, url, tries=4):
    for attempt in range(tries):
        try:
            async with session.get(url, timeout=20) as r:
                if r.status == 404:
                    return None                # don't retry a hard 404
                if r.status < 400:
                    return await r.text()
        except Exception:
            pass
        await asyncio.sleep(2 ** attempt + random.random())  # backoff + jitter
    await dead_letter(url)                     # give up after the cap
    return None

重複排除: 同じページを2回クロールしない

スケールでの発見は常に重複を生み出します。3つのパスから到達可能な同じ製品、1ページを10ページのように見せるトラッキングパラメータ。URLをキューに入れる前に正規化し、見たセット(Redisセット、またはセットが数億に達したらブルームフィルター)を保持します。

from urllib.parse import urlsplit, urlunsplit, parse_qsl, urlencode

seen = set()

def normalize(url):
    s = urlsplit(url.lower())
    q = [(k, v) for k, v in parse_qsl(s.query) if not k.startswith("utm_")]
    return urlunsplit((s.scheme, s.netloc, s.path.rstrip("/"), urlencode(sorted(q)), ""))

def enqueue(url):
    key = normalize(url)
    if key not in seen:
        seen.add(key)
        queue.put_nowait(key)

ストレージとデータ品質

スケール特有の2つの習慣: ストレージがボトルネックにならないようにバッチで書き込み、セレクタが変更されたときに再クロールせずに再解析できるように生のデータと解析済みデータを分けます。そして、多くのチームが見逃す監視の半分を追加します—データ品質チェックです。実行は100%のHTTP成功を報告できますが、レイアウトがずれてセレクタが何も一致しなくなった場合、ゴミを生成することがあります。必須フィールドが空でないこと、値が妥当であることを確認します。

def validate(row):
    assert row.get("title"), "empty title — selector may have drifted"
    price = row.get("price")
    assert isinstance(price, (int, float)) and 0 < price < 1_000_000, "bad price"
    return row

# fail loud on page 5,000, not silently after 5,000,000 empty rows

プロキシ、アンチボット、レンダリングを1つのAPIにオフロードする

必要な場合にのみレンダリング

ヘッドレスブラウザはパイプラインで最も高価な操作です—CPU、メモリ、ページあたりの秒数であり、100万ページではすべてを支配します。多くのサイトはまだ初期のHTMLやJSONエンドポイントにデータを送信しています。プレーンなフェッチとパーサーは桁違いに安価です。まず安価なパスを試し、フィールドが存在することを確認し、必要なページのみレンダリングにエスカレートします。レンダリングコスト対HTTPの内訳でそのギャップに実際の数字を示しています。

難しいレイヤーを構築するか購入するか

上記のすべては構築可能ですが、正直な質問はどの部分がエンジニアリング時間に値するかです。データモデル、解析ロジック、品質チェック、ストレージスキーマはプロジェクトに特有のものであり、あなたがそれをうまく構築できます。プロキシプール、アンチボット処理、ヘッドレスレンダーフリート、リトライとデリバリーキューは一般的なインフラであり、構築するのに高価で、ターゲットが進化するにつれて健康を保つのが大変です。それがScraper APIが座るラインです: 誰にとっても同じ部分をレンタルします。構築対購入の内訳でメンテナンスタックスを詳しく説明し、プロダクションソフトウェアとしてスクレイパーを実行することで全体を観測可能に保つことをカバーしています。

よくある質問

数百万ページをどのようにスクレイプしますか?

より大きなループではなく、同時実行性で。URLの発見と取得の間にキューを置き、非同期または分散ワーカーのプールでそれを消化し、IPを回転させてブロックを回避し、一時的な失敗をバックオフでリトライし、URLを重複排除し、ストレージにバッチで書き込みます。順次の100万ページの実行は数日かかりますが、同じジョブを並行ワーカープールで実行すると数時間で完了します。

大規模ウェブスクレイピングに最適なアーキテクチャは何ですか?

キューとワーカーパイプライン: プロデューサーがURLをキューに列挙し(Redis、RabbitMQ、またはKafka)、ワーカーが回転プロキシレイヤーを通じて同時に取得し、JavaScriptが必要なページのみをレンダリングし、失敗をデッドレターキューにリトライし、見たセットで重複を排除し、生データと解析済みデータを別々に保存します。データ品質チェックで監視をラップし、ドリフトが早期に表面化するようにします。

並行してどれだけのリクエストを実行できますか?

ターゲットとあなたの出口に依存し、固定の数ではありません。スクレイピングはIOバウンドであるため、控えめなボックスでも多くのインフライトリクエストを保持できます。約50の同時実行から始め、エラー率を見て、それが上昇するまで上げます。それがあなたの上限です。1台のマシンの限界を超えて、単一ノードを強く押すのではなく、分散ワーカーを追加します。

スケールで失敗をどのように処理しますか?

失敗を想定します—100万リクエストでは1%のエラー率でも10,000ページが失敗します。一時的なエラー(タイムアウト、503)を指数バックオフとジッターでリトライし、試行回数を制限し、持続的な失敗をデッドレターキューに移動して実行をブロックしないようにします。ハード404はリトライしません。収集順序をランダム化することで、同じページで毎回失敗しないようにします。

スケールは主に構築するのが楽しくない部分です: IPの回転、アンチボット、ヘッドレスレンダリング、キュー、リトライ。データに特有の部分—モデル、パーサー、品質チェック—を所有し、誰にとっても同じグラインドである一般的なインフラをレンタルします。キューと同時実行性を最初に正しく設定し、それ以外はその同時実行性が100万の実際のページと接触しても生き残るようにします。

フェッチレイヤーを回転する住宅IPで強化する