どんなスクレイパーもいずれ途中で止まる。コンテナが再スケジュールされたり、デプロイでワーカーが再起動されたり、マシンのメモリが尽きたり、誰かがCtrl+Cを押したりする。次に何が起きるかで、データが信頼できるものであり続けるかどうかが決まる。単純に最初からやり直すスクレイパーは、すべてを再取得し、そのために設計されていない限り、すべてを再度書き込んでしまう。
冪等なスクレイパーとは、同じ作業を2回実行しても1回実行したときと同じ結果になるスクレイパーのことだ。どんな瞬間にkillされても、再起動すれば、各レコードにつき正確に1つの正しいコピーで終わる。本稿では、それを実現する4つの設計判断を、検証済みのコードと、クロール中に10回killした結果とともに示す。
要点
- プロセスはどの2行の間でも死にうると想定する。再起動時に繰り返されるのは、処理中だったページ1件だけになるように設計する。
- すべてのレコードに、取得元やSKUなど「それが何であるか」から導出したIDを付与し、「いつ取得したか」からは決して導出しない。そうすれば、書き込みの繰り返しは重複ではなく更新になる。
- クロールのフロンティアは永続ストレージに保持し、URLを完了とマークする処理は、その結果を保存するトランザクションと同一にする。
- ロックではなくリースを使い、クラッシュしたワーカーが確保していた作業が自動的に再び利用可能になるようにする。
- 500個の商品を対象に、各実行で10回killするテストを行ったところ、素朴なスクレイパーは1,286行から2,165行、最大で4.3倍のリクエスト数で終わったのに対し、冪等なスクレイパーは正確に500行の正しいデータと、最大3回の重複リクエストで終わった。
再起動が重複を生む理由
スクレイパーの最初のよくあるバージョンはこうだ。1ページ目から始めてページネーションをたどり、商品ごとに行を挿入し、逐次コミットしていく。これは中断されるまでは完璧に動く。再起動すると、どこまで進んだかの記憶がないため最初からやり直し、すでに保存済みの商品がすべて、新しい自動採番IDで再度挿入される。何度かkillすると、テーブルにはほとんどの商品について複数のコピーが残り、どれが最新なのかを示す情報は何もない。
後から重複を排除することも可能で、エンティティ解決がそれを扱っているが、それは症状への対処にすぎない。重複は精度を損なう前にコストもかかる。再取得されたページはすべて、2回分の帯域コストであり、それはそのままクリーンレコードあたりのコストに跳ね返る。
スクレイパーを冪等にする4つの設計判断
1. レコードIDをデータから導出する
レコードのIDは、そのレコードが何であるか、すなわち取得元と、SKUやリスティングID、正規化されたURLといった自然キーから導出すべきだ。それらをハッシュ化すれば、同じ商品は、どの実行でもどのマシンでも常に同じIDになる。再度の書き込みはアップサートになる。レコードが既に存在し変化がなければ何も起きず、変化していれば更新され、変更時刻が記録される。
2. フロンティアを永続化する
訪問すべきURLのリストと、どれが完了したかという情報は、メモリではなくデータベースに置くべきだ。既に知られているURLの追加は何もしない操作でなければならず、既に処理済みの一覧ページでリンクを再発見しても害はない。
3. 結果と進捗を一緒にコミットする
危険な瞬間は、結果の保存とURLが完了したという記録との間にある。プロセスが一方の後もう一方の前に死ねば、再起動は結果を失うか、結果を繰り返すかのどちらかになる。両方を1つのデータベーストランザクションで行えば、この隙間はなくなる。レコードが保存されURLが完了とマークされるか、どちらも起きずURLが単に再度取得されるか、そのどちらかになる。
4. ロックではなくリースで作業を管理する
ワーカーは限られた時間だけURLを確保する。処理が終われば、URLは完了とマークされる。クラッシュした場合はリースが期限切れになり、別のワーカー、あるいは再起動した同じワーカーがそのURLを引き取る。作業が回復するために、クラッシュに誰かが気づく必要はない。
コード
以下のモジュールは、この4つすべてをSQLiteとrequestsで実装したものだ。SQLiteを使うことで例を自己完結させているが、同じ設計はPostgreSQLや、トランザクションとアップサートを持つ任意のデータベースにもそのまま適用できる。
import hashlib
import json
import re
import sqlite3
import time
import requests
SCHEMA = """
CREATE TABLE IF NOT EXISTS frontier (
url TEXT PRIMARY KEY,
status TEXT NOT NULL DEFAULT 'pending', -- pending, leased, done
leased_until REAL NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS records (
record_id TEXT PRIMARY KEY, -- derived from the source and its natural key
url TEXT NOT NULL,
content_hash TEXT NOT NULL,
data TEXT NOT NULL,
first_seen REAL NOT NULL,
last_changed REAL NOT NULL
);
"""
def open_db(path):
db = sqlite3.connect(path, isolation_level=None) # explicit transactions below
db.execute("PRAGMA journal_mode=WAL")
db.executescript(SCHEMA)
return db
def record_id(source, natural_key):
"""The same item always gets the same ID, however many times it is scraped."""
return hashlib.sha256(f"{source}:{natural_key}".encode()).hexdigest()[:16]
def enqueue(db, urls):
"""Adding a URL that is already known is a no-op, so re-discovering links is harmless."""
db.executemany("INSERT OR IGNORE INTO frontier (url) VALUES (?)", [(u,) for u in urls])
def lease(db, lease_seconds=60):
"""Claim one URL. A lease left behind by a crashed worker expires and the URL is retried."""
now = time.time()
row = db.execute(
"UPDATE frontier SET status = 'leased', leased_until = ? WHERE url = ("
" SELECT url FROM frontier WHERE status = 'pending' OR (status = 'leased' AND leased_until < ?) LIMIT 1"
") RETURNING url", (now + lease_seconds, now)).fetchone()
return row[0] if row else None
def complete(db, url, new_urls=(), record=None):
"""Write the result and mark the URL done in one transaction: either both happen or neither does."""
db.execute("BEGIN IMMEDIATE")
try:
enqueue(db, new_urls)
if record:
now = time.time()
payload = json.dumps(record["data"], sort_keys=True)
digest = hashlib.sha256(payload.encode()).hexdigest()
db.execute(
"INSERT INTO records (record_id, url, content_hash, data, first_seen, last_changed) VALUES (?, ?, ?, ?, ?, ?) "
"ON CONFLICT(record_id) DO UPDATE SET data = excluded.data, content_hash = excluded.content_hash, "
"last_changed = excluded.last_changed WHERE records.content_hash != excluded.content_hash",
(record["id"], url, digest, payload, now, now))
db.execute("UPDATE frontier SET status = 'done' WHERE url = ?", (url,))
db.execute("COMMIT")
except Exception:
db.execute("ROLLBACK")
raise
def crawl(db, base, session=None, lease_seconds=60):
session = session or requests.Session()
enqueue(db, [f"{base}/list/0"])
while True:
url = lease(db, lease_seconds)
if url is None:
if db.execute("SELECT 1 FROM frontier WHERE status != 'done' LIMIT 1").fetchone():
time.sleep(1) # a lease is still held, perhaps by a worker that crashed; wait for it to expire
continue
return
html = session.get(url, timeout=30).text
if "/list/" in url:
links = [base + href for href in re.findall(r'href="(/(?:product|list)/\d+)"', html)]
complete(db, url, new_urls=links)
else:
sku = re.search(r'class="sku">([^<]+)<', html).group(1)
price = re.search(r'class="price">([^<]+)<', html).group(1)
complete(db, url, record={"id": record_id("example-shop", sku), "data": {"sku": sku, "price": price}})
crawl関数はテストサイト固有のもので、25の一覧ページが500の商品ページにリンクしている。それより上の部分はすべて再利用可能だ。見落としやすい点が2つある。
- アップサートは内容が変化したときだけ書き込む。 競合更新の
WHERE句は、変化していないレコードをそのままにしておくため、last_changedはその名の通りの意味を持ち、変更検知の駆動に使える。 - ループは空のキューで止まらない。 最初のバージョンは、リースできるURLがなくなった時点で終了していたため、クラッシュによってリースが残っている場合にページが未完了のまま残ってしまっていた。現在のループは、すべてのURLが完了するまで待機する。
10回killしてみる
25の一覧ページと500の商品ページを持つローカルのテストサイトを用意した。完全なクロールには正確に525回のリクエストが必要になる。各実行では、スクレイパーを起動し、0.2秒から1秒の間のランダムな瞬間にSIGKILLでkillし、それを10回繰り返してから最後に完了させた。上記の冪等なスクレイパーと、先に説明した素朴なスクレイパーを、同じkillのタイミングでそれぞれ5回ずつ実行し、冪等版のリースは3秒とした。
| Run | 素朴な方式: 保存行数 | 素朴な方式: リクエスト数 | 冪等な方式: 保存行数 | 冪等な方式: リクエスト数 |
|---|---|---|---|---|
| 1 | 2,165 | 2,283 | 500 | 525 |
| 2 | 1,875 | 1,977 | 500 | 526 |
| 3 | 1,862 | 1,967 | 500 | 526 |
| 4 | 1,455 | 1,538 | 500 | 526 |
| 5 | 1,286 | 1,361 | 500 | 528 |
どちらの方式も、最後の実行を完了させたため、最終的には500個すべての商品を保存した。違いはそれ以外のすべてにある。素朴なスクレイパーは商品1件あたり2.6行から4.3行を保存し、必要なリクエスト数の2.6倍から4.3倍を発行した。冪等な方式は商品1件あたり正確に1行を保存し、すべての値が取得元ページと一致しており、10回のクラッシュを通じて重複リクエストは最大3回にとどまった。これは、ページが処理中だったときにkillが発生した回数に対応する。また、クラッシュしたワーカーのリースが期限切れになるのを待つこともあったにもかかわらず、すべての実行で所要時間はより短く、9.2秒から10.8秒に対し6.8秒から8.3秒だった。
データベース以外の副作用
スクレイパーが2回以上行うのはレコードの書き込みだけではない。同じ原則は、あらゆる副作用に当てはまる。
- ファイル。 ダウンロードした画像や文書は、内容のハッシュ値やレコードIDによって命名し、再ダウンロードがコピーを追加するのではなく上書きするようにする。
- メッセージとWebhook。 レコードとそのバージョンから導出した冪等性キーを送信し、消費側が既に処理済みのメッセージを無視できるようにする。
- カウンターと集計値。 取得のたびにインクリメントするのではなく、保存済みのレコードから再計算する。そうしないと、再起動のたびに値が膨らんでしまう。
- 生のレスポンス。 生のレスポンスをアーカイブする場合、再取得は2つ目のキャプチャを追加するだけであり、抽出されたレコードが冪等である限り、無害でむしろ有用だ。
実践的なルール
- 自然キーは慎重に選ぶ。 それは取得元において安定していなければならない。ページ上の位置やトラッキングパラメータ付きのURLではなく、SKUやリスティングIDを使う。
- まずリトライを安全にし、その後頻度を上げる。 書き込みが冪等であれば、リトライとバックオフをデータを汚さずに気前よく行える。
- リースの長さは最も遅いページに合わせる。 遅い取得よりも短いリースでは、2つのワーカーが同じURLを処理することになる。ここでは害はないが、帯域の無駄になる。
- 繰り返し率を監視する。 保存レコードあたりのリクエスト数は、低コストな健全性指標だ。これが上昇すればクラッシュやリースの問題を意味し、スクレイピングパイプラインの監視で扱う他の指標と並べて管理すべきものだ。
結論
安全に再起動できないスクレイパーは、いずれ自らのデータを破損させる。しかもそれは静かに起きる。解決策は、より注意深い運用ではなく、再起動を退屈なものにする設計だ。データから導出されたID、永続的なフロンティア、一緒にコミットされる結果と進捗、そして期限切れになるリース。
今回のテストでは、この4つの設計判断によって、10回のクラッシュが、最大4.3倍の行数とリクエスト数をもたらす状態から、正確に正しい500件のレコードと3回の重複リクエストへと変わった。