Web Scraping APIは収集における最も難しい部分、プロキシ、レンダリング、リトライ、ブロックへの対処を取り除いてくれる。取り除いてくれないのは、レスポンスが到着した後に起きるすべてのことであり、実際のところ多くのパイプラインはそこで失敗する。
失敗は静かに起きる。誰も重複排除しなかったリトライによる重複行。誰かがダッシュボードで数値にキャストした"$16.99"のような文字列だらけの価格列。3週間前にサイトのリデザインでフィールドがnullになっていたこと。もう一度取得料金を払わなければ直せないエクストラクタのバグ。
このガイドはロード側、つまりスクレイピングAPIの出力をSQLデータベースに、量をこなしながら確実に、かつ保存内容の説明やリプレイができる状態を保ちながら投入することを扱う。例にはPostgreSQLとShifter Web Scraping APIを使うが、パターンはどのSQLデータベースにも通用する。
実際に得られるレスポンスから始める
Shifter Web Scraping APIでは2つのレスポンス形式から選べる。
生のHTMLで、自分側でパースする方法。あるいはextract_rulesを渡してCSSセレクタをフィールドにマップする構造化JSONだ。
curl "https://scrape.shifter.io/v1?api_key=YOUR_API_KEY&url=https://shop.example.com/p/42&render_js=1&extract_rules=%7B%22title%22%3A%7B%22selector%22%3A%22h1%22%2C%22output%22%3A%22text%22%7D%2C%22price%22%3A%7B%22selector%22%3A%22.price%22%2C%22output%22%3A%22text%22%7D%7D"
# {"title": "Example Product", "price": "$19.99"}
この出力形式には、スキーマに影響する2つの特性がある。セレクタが何にもマッチしないフィールドはリクエストが失敗するのではなくnullとして返ってくるため、要素が存在しないことと壊れたセレクタとがレスポンス上では同じに見える。またテキスト出力は表示用の文字列なので、価格、評価、日付は人間向けにフォーマットされた形で届き、データベース向けに型付けされているわけではない。結果ページのリスト抽出を含むルールの構文はextraction rules docsにある。
1つのテーブルではなく3つの層
本番環境と接触して生き残る設計は、受け取ったものと結論づけたものを分離する。
| 層 | 内容 | 存在理由 |
|---|---|---|
| Raw landing | フェッチメタデータ付きで受信したままのすべての成功レスポンス | 再フェッチせずに抽出をリプレイできる |
| Typed observations | パース済み、型付け済み、検証済みの値とパースステータス | アナリストやアプリケーションがクエリする対象 |
| Current state | observationsから導出された、エンティティごとの最新値 | 製品やダッシュボード向けの高速な読み取り |
raw層はチームが省略しがちで、後で後悔する層だ。クレジットは成功したリクエストに対して消費されるため、再度フェッチしなければ直せないパースのバグは、クロール全体のコストを二重にする。まずレスポンスを着地させ、次にパースすれば、パーサーの修正はリプレイクエリで済む。
landingテーブル
CREATE TABLE scrape_raw (
job_id text PRIMARY KEY,
source_url text NOT NULL,
market text NOT NULL,
fetched_at timestamptz NOT NULL,
http_status smallint NOT NULL,
body jsonb NOT NULL,
body_hash text NOT NULL,
extractor_version text NOT NULL
);
CREATE INDEX scrape_raw_url_time ON scrape_raw (source_url, fetched_at DESC);
ここにはいくつかの意図的な選択がある。
job_idは冪等性キーであり、URLとスケジューリングウィンドウからリクエスト前に計算される。そのため、同一の論理ジョブのリトライは2行目を作るのではなく同じキーに着地する。extractor_versionはどのエクストラクションルールがそのbodyを生成したかを記録し、これによって後でサイトの変更とルールの変更を区別できるようになる。marketはどこで観測されたかを記録する。同一URLが国ごとに異なるコンテンツを返すことがあるからだ。そしてAPIキーは、保存するリクエストメタデータのどこにも保存されない。データベース内の認証情報はすべてのバックアップ内の認証情報になるからだ。
冪等にロードする
import hashlib
import json
import os
import psycopg
import requests
from psycopg.types.json import Jsonb
API = "https://scrape.shifter.io/v1"
RULES = {
"title": {"selector": "h1", "output": "text"},
"price": {"selector": ".price", "output": "text"},
}
EXTRACTOR_VERSION = "product-v3"
INSERT_RAW = """
INSERT INTO scrape_raw
(job_id, source_url, market, fetched_at, http_status, body, body_hash, extractor_version)
VALUES (%s, %s, %s, now(), %s, %s, %s, %s)
ON CONFLICT (job_id) DO NOTHING
"""
def job_id(url: str, market: str, window: str) -> str:
return hashlib.sha256(f"{url}|{market}|{window}".encode()).hexdigest()
def fetch(url: str, market: str) -> requests.Response:
params = {
"api_key": os.environ["SHIFTER_API_KEY"],
"url": url,
"render_js": 1,
"country": market,
"extract_rules": json.dumps(RULES),
}
return requests.get(API, params=params, timeout=90)
def land(conn: psycopg.Connection, url: str, market: str, window: str) -> None:
resp = fetch(url, market)
if resp.status_code != 200:
raise RuntimeError(f"{resp.status_code} for {url}")
body = resp.text
with conn.cursor() as cur:
cur.execute(
INSERT_RAW,
(
job_id(url, market, window),
url,
market,
resp.status_code,
Jsonb(json.loads(body)),
hashlib.sha256(body.encode()).hexdigest(),
EXTRACTOR_VERSION,
),
)
ON CONFLICT (job_id) DO NOTHINGによって、このロードはリトライしても安全になる。失敗の原因がネットワークであれ、ワーカーであれ、データベースであれ、ジョブを再度実行しても重複は生まれない。MySQLでの相当する処理は、ユニークキーとINSERT IGNOREまたはON DUPLICATE KEY UPDATEだ。
スループット: フェッチとロードを分離する
フェッチ側にはプランの同時実行数上限によるハードな上限があり、それを超えるリクエストは429を返す。データベースにも独自の上限があり、通常はスクレイピングしたページごとに接続とトランザクションを開くチームが最初にぶつかる上限だ。
両者の間にキューを置く。同時実行数上限に合わせてサイズ調整されたフェッチワーカーが、レスポンスをキューに書き込む。少数のローダーワーカーがバッチでキューを消費する。安定した流量であれば、psycopgのexecutemanyを数百行のバッチで使えば十分だ。大規模なバックフィルの場合は、バッチをCOPYでunloggedのステージングテーブルに入れ、1つのステートメントでマージする。
INSERT INTO scrape_raw
SELECT * FROM scrape_raw_staging
ON CONFLICT (job_id) DO NOTHING;
これによって何千回ものラウンドトリップが1回になり、冪等性の保証もそのまま維持される。
長時間かかるレンダリングについては、APIは非同期で配信できる。webhook=<URL>を渡せば、準備できたレスポンスがそのエンドポイントにPOSTされる。そのレシーバーもjob_idで冪等にしておくこと。どちらの側からでもHTTP配信はリトライされる可能性があるためだ。
APIエラーをパイプラインの挙動にマッピングする
ステータスコードはすべてがリトライ対象というわけではなく、それらを一律に扱うローダーは、壊れた設定にスパムを送るか、一時的な失敗を諦めてしまうかのどちらかになる。完全な表はerrors and limitsにある。
| ステータス | パイプラインの挙動 |
|---|---|
408、422、500 | 指数バックオフでリトライ |
429 | バックオフしてワーカーの同時実行数を減らす |
400、401、403 | 設定エラー: デッドレターキューに送りアラートを出す、決してリトライしない |
509 | クレジット枯渇: フェッチ段階を停止しアラートを出す、リトライしても意味がない |
失敗したリクエストとターゲットの4xxや5xxレスポンスには課金されず、APIはすでに返却前に一時的な失敗を最大3回リトライしているため、自分側のリトライはクレジットではなく時間のコストになる。それでも時間はかかるため、バックオフが重要になる。
observationsを型付けする
これが表示用文字列がデータになる場所であり、大半のサイレントなエラーが持ち込まれる場所だ。
CREATE TABLE price_observation (
source_url text NOT NULL,
market text NOT NULL,
observed_at timestamptz NOT NULL,
price_amount numeric(12,2),
currency char(3),
raw_price text,
parse_status text NOT NULL,
job_id text NOT NULL REFERENCES scrape_raw (job_id),
PRIMARY KEY (source_url, market, observed_at)
);
3つのルールがこれを正直に保つ。
パース済みの値の隣にraw文字列を保持する。 raw_priceがあるからこそ、再フェッチせずに疑わしい数値を監査できる。
グローバルにではなく、市場ごとにパースする。 "1.299,00"と"1,299.00"は表記慣習が違うだけで同じ価格であり、通貨記号は通貨そのものではない。$はストアフロントによって米ドル、カナダドル、オーストラリアドルのいずれかを意味する。記号と市場の両方からISOコードを解決すること。
値がnullである理由を記録する。 missing、unparseable、okというparse_statusは、「ページに価格がなかった」ことと「パーサーが失敗した」ことを区別する。これはAPIのnull単体では区別できないことだ。
現在の状態については、維持するのではなく導出する。PostgreSQLでは:
CREATE VIEW price_current AS
SELECT DISTINCT ON (source_url, market) *
FROM price_observation
WHERE parse_status = 'ok'
ORDER BY source_url, market, observed_at DESC;
導出されたビューは、それが要約する履歴と食い違うことがない。
ユーザーより先にスキーマドリフトを検知する
サイトはマークアップを変更するが、変更されたセレクタはエラーを出さない。nullを返し、リクエストは成功し、クレジットが消費され、その行は一見有効に見える状態で着地する。
防御策は、フィールドごと、ソースごと、エクストラクタバージョンごとのnull率モニターだ。各フィールドのmissingパースステータスの割合をローリングウィンドウで計算し、ベースラインから急激に動いたらアラートを出す。価格フィールドが欠損2%から一晩で60%欠損に変化したとすれば、それはリデザインであり、同日中に検知できるかどうかが、修正済みセレクタと3週間の使えない履歴データとの違いになる。
修正したら、extractor_versionを上げ、影響を受けたraw行を新しいパーサーで再処理する。それがrawレスポンスを着地させることの見返りだ。
保持期間とパーティショニング
raw landingテーブルは最も速く成長し、最も読まれることが少ない。フェッチ日でパーティション分割し、現実的なリプレイウィンドウをカバーするのに十分な期間だけ保持し、行を削除するのではなく古いパーティションをドロップする。observationsテーブルは履歴記録であり、通常は同様にパーティション分割された、より長い保持期間を確保すべきだ。
何を監視すべきか
| メトリック | 何を検知するか |
|---|---|
| 成功したフェッチ数とロードされた行数の比較 | APIとデータベース間でのローダーの取りこぼし |
重複したjob_idのコンフリクト | リトライの嵐やスケジューリングの重複 |
| フィールドおよびエクストラクタバージョンごとのnull率 | マークアップの変化や壊れたセレクタ |
| キューの深さとロードの遅延 | フェッチ段階に遅れているローダー |
消費されたクレジットとokにパースされた行数の比較 | 使えないレスポンスに費やされたお金 |
最後のメトリックが重要なコストの見方だ。リクエストあたりのクレジットではなく、使える行あたりのクレジットである。API自体の使用状況とエラー率は、パネル内のWeb Scraping APIの項で確認できる。
FAQ
raw層にはHTMLと抽出済みJSONのどちらを保存すべきか?
抽出済みJSONの方がはるかに小さく、通常は十分だ。HTMLの保存は、抽出ロジックを頻繁に変更すると見込まれるソースに限定し、そのテーブルには短い保持期間を設定すること。
JSONBは直接クエリするのに十分か?
探索的な用途であれば十分だ。プロダクトやダッシュボードが依存する何かについては、フィールドを型付きカラムに昇格させ、データベースが型を強制し通常のインデックスを使えるようにすること。
変化していないページへの課金をどう避けるか?
成功したリクエストはすべてクレジットを消費するため、節約はより少なく書き込むことからではなく、より少なくフェッチすることから生まれる必要がある。どの詳細ページをそもそもフェッチする必要があるかを判断するために、一覧ページやサイトマップページのような安価なシグナルを使うこと。
これはAmazon APIでも機能するか?
機能する。そのレスポンスはすでに構造化されたJSONであるため、エクストラクションルールは不要になるが、価格は依然として表示用文字列として届くため、landing、typing、driftのパターンはそのまま適用できる。
結論
スクレイピングAPIは収集を解決する。データベースの設計が、収集したものが信頼できる状態を保つかどうかを決める。すべての成功レスポンスを冪等性キーとエクストラクタバージョンとともに着地させ、raw文字列を保持しながら市場ごとに値を型付けし、値がnullである理由を記録し、現在の状態を維持するのではなく導出し、フィールドごとのnull率を監視してマークアップの変化が同日中に表面化するようにすること。
Amazonに特化した内容については、the best web scraping APIs for Amazon monitoringで選択肢を比較しており、同じパイプラインの不動産版はhow real estate companies use web scraping APIsにある。製品はWeb Scraping APIページにあり、プランはpricing pageにある。