あらゆる分散型クローラーは、いずれ同じ悪夢のような午後を迎える。あるサイトが巨大なページを配信し始めたせいで、パーサーの処理が遅くなる。フェッチャーは何も言われないため、フルスピードで取得を続ける。両者の間のキューは数百万件単位で膨れ上がり、メモリ使用量が上昇し、あるノードが落ち、その作業は他のノードに再キューイングされ、再試行がそれらも道連れにする。その間、そもそもの原因となったそのサイトは、かつてないほどのトラフィックを受け続けている。
これはどの単一コンポーネントのバグでもない。コンポーネント間にフロー制御が存在しないことが原因だ。クローラーとは速度がまったく異なるステージのパイプラインであり、遅いステージが速いステージに「待て」と伝える手段がなければ、速いステージがすべてを台無しにするまで勝ち続ける。
このエッセイでは、そのフィードバックをどう組み込むかを扱う。シグナルはどこから来て、どこに伝わるべきか、そしてクローラーを崩壊させずに緩やかに劣化させるための、わずかなルール群だ。
クローラーはパイプラインであり、各ステージは速度について一致しない
細部を取り除けば、たいていのクローラーは4つのステージを持つ。
| ステージ | 何によって制限されるか | 過負荷時にどう壊れるか |
|---|---|---|
| フロンティアとスケジューラ | ほぼ何もない、URLを選ぶだけ | 下流が吸収できるより速く作業を発行する |
| フェッチャー | 対象サイト、プロキシ容量、プランの上限 | タイムアウト、スロットリング、ブロック |
| パースとレンダリング | CPU、メモリ、ヘッドレスブラウザのスロット | レイテンシの増大、その後メモリ枯渇 |
| 重複排除とストレージ | データベースの書き込み容量 | 書き込みの遅延、ロック競合、バックログ |
フロンティアはほぼ無料で動くため、常に最速のステージになる。他の3つのステージはそれぞれ異なる要因で制限されており、その制限は移り変わる。あるサイトが遅くなったり、あるページ種別が重くなったり、データベースがコンパクションを始めたりする。フロー制御とは、現時点で最も遅いステージがすべてのペースを決めることを意味する。
その代替は、どのステージが最初にメモリを使い果たすかによってペースが決まることだ。
ルール1: すべてのキューには上限を設ける
2つのステージ間の無制限のキューは、バッファではない。どんな量の余剰作業でも吸収するという、どのマシンも守れない約束だ。それは問題を隠しもする。上流のステージはエンキューのたびに成功を目にするが、キューは静かに膨れ上がり、誰かが気づいたときには最も古いアイテムは何時間も待たされている。
上限付きキューは同じ状況をシグナルに変える。キューが満杯になると、プロデューサーはブロックされるか、明示的な拒否を受け取る。減速は上流に一段ずつ伝播し、最終的にはフロンティアに到達して、フロンティアはしばらくの間URLの配布を単に止める。何も失われず、何も爆発しない。
単一プロセス内であれば、これはキューのコンストラクタへの一つの引数にすぎない。
import asyncio
async def fetcher(frontier: asyncio.Queue, parsed: asyncio.Queue, fetch):
while True:
url = await frontier.get()
try:
page = await fetch(url)
# Blocks when the parse queue is full: a slow parser slows fetching.
await parsed.put(page)
finally:
frontier.task_done()
async def parser(parsed: asyncio.Queue, store):
while True:
page = await parsed.get()
try:
await store(page)
finally:
parsed.task_done()
async def run(urls, fetch, store, fetchers=16, parsers=4):
frontier = asyncio.Queue(maxsize=1000)
parsed = asyncio.Queue(maxsize=200)
workers = [asyncio.create_task(fetcher(frontier, parsed, fetch)) for _ in range(fetchers)]
workers += [asyncio.create_task(parser(parsed, store)) for _ in range(parsers)]
for url in urls:
await frontier.put(url) # blocks when the frontier is full
await frontier.join()
await parsed.join()
for w in workers:
w.cancel()
複数マシンにまたがっても原則は同じで、変わるのは仕組みだけだ。長さ制限と拒否ポリシーを持つブローカーキュー、実際に対処するコンシューマーラグを持つストリーム、あるいは空き容量のあるワーカーにしかURLを渡さないデータベース上のフロンティアなどだ。サイズの上限は、手元にあるメモリからではなく、許容できるレイテンシから決める。パースステージが1秒あたり200ページを処理でき、10秒間のキューイングを許容するなら、キューは約2,000アイテムを保持する。それより大きくしても、問題に気づく瞬間を先延ばしにするだけだ。
ルール2: プッシュではなくプルにする
最もクリーンなバックプレッシャーは、無償で手に入る種類のものだ。コーディネーターがワーカーに作業を割り当てるのではなく、ワーカーが余力のあるときに作業を要求するようにすれば、遅いワーカーは単に要求する頻度を下げるだけになる。システムはワーカーが処理できる以上のものを送りつけることができない。なぜなら、それは決して要求されていないからだ。
作業をプッシュすると、コーディネーターは推測せざるを得なくなる。各ワーカーの負荷を追跡しなければならず、負荷が変化するまさにその瞬間に推測を誤る。プル型の設計は、ワーカーがキューからアイテムをリースするにせよ、上流に対して明示的なクレジットを与えるにせよ、その判断を答えを知っている唯一のコンポーネントに委ねる。
リースには有効期限が必要だ。作業を保持したまま死んだワーカーが、それを永遠に保持し続けてはならない。そのため、リースされたアイテムはタイムアウト後にキューに戻される。そのタイムアウトは平均ではなく、正当に最も遅いフェッチから設定する。そうしなければ、健全だが遅い作業が二重に配布されてしまう。
ルール3: グローバルなキューではなく、ホストごとにキューを持つ
単一のグローバルなフェッチキューには、クローラーが常に陥る故障モードがある。ヘッドオブラインブロッキングだ。次の1,000件のURLがすべて一つの遅い、あるいはスロットリング中のサイトに属している場合、すべてのフェッチャーはそのサイトを待つことになり、その背後には速く健全なサイト向けの作業が滞留する。
フェッチステージをホストごと、あるいは対象がレート制限をかける単位ごとに分割する。各ホストは自分自身のキューと自分自身の並行数の上限を持ち、フェッチャーは現在余裕のあるホストから選び取る。すると、苦戦している一つのサイトは自分自身だけを遅くする。
これはまた、礼儀正しさが宿る場所でもある。ホストごとの上限は、自分に対するフロー制御であると同時に、サイトに対する自制でもある。その上限をそもそもどう設定するか、そしてサイト自身のシグナルがそれについて何を教えてくれるかについては、レート制限とリクエストスロットリングで扱っている。
ルール4: ホストごとの上限を適応的にする
固定されたホストごとの並行数は、両方向で誤っている。低すぎれば、もっと受け入れられるサイトの容量を無駄にする。高すぎれば、押し返し始めたサイトに圧力をかけ続けることになり、それが一時的なスロットリングをブロックに変える。
十分に検証された答えは、TCPが使っているものと同じだ。加法的増加、乗法的減少。成功のたびに上限を少しずつ引き上げる。スロットリングやタイムアウトのたびに半分に減らす。上限はそのサイトが許容する水準のわずかに下に落ち着き、サイトが変化すればそれに応じて動く。
import asyncio
class HostLimiter:
"""Adaptive concurrency for one host: additive increase, multiplicative decrease."""
def __init__(self, start=4, floor=1, ceiling=32):
self.limit = start
self.floor = floor
self.ceiling = ceiling
self.in_flight = 0
self._cond = asyncio.Condition()
async def acquire(self):
async with self._cond:
await self._cond.wait_for(lambda: self.in_flight < int(self.limit))
self.in_flight += 1
async def release(self, outcome):
async with self._cond:
self.in_flight -= 1
if outcome == "ok":
self.limit = min(self.ceiling, self.limit + 1 / max(1, int(self.limit)))
elif outcome in ("throttled", "timeout"):
self.limit = max(self.floor, self.limit / 2)
self._cond.notify_all()
増加は意図的に緩やかで、成功が一巡するごとにおよそ1スロット追加される程度であり、減少は意図的に急激だ。この非対称性が要点だ。サイトの許容量を超過することのコストは、下回ることのコストをはるかに上回る。自分で設定したホストごとのハード上限は必ず維持し、適応的な上限がそれを妥当だと考える範囲を超えて動くことのないようにする。
ルール5: どのシグナルがどこから来たかを把握する
「減速せよ」がすべて同じ意味を持つわけではなく、それらを同じように扱うと、圧力が誤った場所に送られてしまう。
| シグナル | 発生源 | 何を減速させるべきか |
|---|---|---|
サイトからの429、レイテンシの上昇、チャレンジページ | 対象サイト | そのホストのリミッターのみ |
| パースキューの満杯、ストレージの遅延 | 自分自身のパイプライン | フェッチ全体、続いてフロンティア |
| スクレイピングAPIの並行数上限 | 自分のプラン | そのAPIへの合計インフライトリクエスト数 |
| クォータまたはクレジットの枯渇 | 自分のプラン | サイクルがリセットされるか、追加補充するまでのすべて |
例えば、ShifterのWeb Scraping APIは、プランの並行数上限を超えると429 Too Many Requestsを返し、クレジットが枯渇すると509を返す。前者は対象サイトについてではなく、自分自身についてのフロー制御シグナルであり、そのため、どのホストごとのリミッターにも属さず、そのAPIへの呼び出しに対するグローバルなリミッターに属する。後者は停止条件だ。どちらかをサイト自身の429と混同すると、クローラーは何の関係もない問題のために健全なサイトをスロットリングしてしまう。
リトライは負荷である
クローラーが自らのバックプレッシャーを台無しにする最も一般的な方法は、リトライを通じてだ。フェッチが失敗し、ワーカーは即座に再試行し、原因が解消していないため再試行も失敗し、同じことをしているすべてのワーカーが、対象サイトあるいは自分自身のステージが最も余裕のないまさにその瞬間にトラフィックを何倍にも増やす。
- リトライは初回試行と同じ受付制御を通過させる。 ホストごとのリミッターをスキップするリトライは、自分自身のフロー制御を回避する抜け道だ。
- ジッターを加えてバックオフする。 そうすれば、一緒に失敗したワーカー群が一緒に再試行しなくなる。
- 各ジョブにリトライ予算を与え、リトライを総リクエスト数に対する割合として追跡する。その割合が上昇しているなら、システムは失敗にキャパシティを費やしている。
- リトライ層を積み重ねない。 Web Scraping APIはすでに、失敗したフェッチ、CAPTCHA、一時的な対象エラーを、異なるプロキシを使って最大3回まで、追加料金なしで再試行してからエラーを返している。その上にさらに積極的なクライアント側のリトライを重ねると、レジリエンスが増すのではなく、試行回数が増えるだけになる。
- 成功しえないものは決してリトライしない。 認証エラーや設定エラーは、毎回同じように失敗する。
追いつけないときは、意図的に切り捨てる
時には、クローラーがこなせる以上の作業を抱えているというのが正直な答えになる。バックプレッシャーはその場合、すべてを均等に遅くするが、それはしばしば最悪の結果になる。すべてのジョブが遅れて完了することになり、重要なジョブも例外ではない。
負荷の切り捨ては、その選択を明示的にする。作業に優先度を与え、キューが閾値を超えて満杯であり続けるなら、フェッチにコストをかける前に、フロンティアの段階で最も優先度の低いアイテムを破棄または延期する。頻繁に変動する価格ページを更新することは、1年間変化していないアーカイブページを再チェックすることより価値が高い。どのページがその予算に値するかはそれ自体が別の問題であり、コストを意識したクロールスケジューリングで扱っている主題だ。
スループットではなく、圧力を測定する
スループットのグラフは、崩壊が起きる直前まで問題なく見える。圧力の高まりを示す指標は別のものだ。
- 各ステージのキューの深さと、最も古いアイテムの経過時間。 深さよりも経過時間の方が重要だ。深いキューでも素早く捌けているなら健全であり、浅くても古いアイテムで満ちているキューは健全ではない。
- 各ホストのインフライトリクエスト数と、現在の適応的な上限。 半減し続ける上限は、サイトが押し返している証拠だ。
- 各ステージの受付拒否とブロックされたput。 これはバックプレッシャーがその役割を果たしている証拠だ。急激な上昇は、現在どのステージがボトルネックになっているかを教えてくれる。
- 各ホストおよび全体でのリトライ比率。
ホストごとの上限が下がり続けることと、ブロック率およびチャレンジ率の上昇が組み合わさっている場合、それは通常、サイトが単に遅くなっているのではなく、自分に対して敵対的になりつつあることを意味する。これらのシグナルをサイトごとの単一の判断に変える方法については、対象サイトのヘルススコアの構築で扱っている。より広範なパイプライン指標については、Webスクレイピングパイプラインの監視にある。
水平スケーリングはフロー制御の必要性を取り除かない
ワーカーを追加すればフェッチステージの上限は上がる。それはどの対象サイトの許容量も、パース容量も、プランの上限も引き上げない。ホストごとの上限と上限付きの下流キューを持たずに、キューの深さに応じてフェッチャーを自動でスケールさせるクローラーは、すべての制約に一斉に向かってスケールしてしまう。Kubernetesで運用しているなら、CPUだけでなく上記のシグナルに基づいてスケールさせるべきであり、その仕組みについてはKubernetesによるレジデンシャルプロキシスクレイピングのスケーリングで扱っている。
同じことはリージョンをまたいでも当てはまる。永続的なキューとリージョン間のバックプレッシャーは、あるリージョンの障害がグローバルな障害に転じるのを防ぐものであり、マルチリージョンパイプラインにおけるレジデンシャルプロキシのフェイルオーバーで説明されている。
結論
フロー制御は、クローラーが大規模になってから追加する最適化ではない。それは、負荷にさらされたクローラーが減速するか崩壊するかを決めるものだ。ルールはわずかだ。すべてのキューに上限を設け、ワーカーにプルさせ、ホストごとに分割し、ホストごとの上限を急激に下げつつ緩やかに上げ、それぞれの減速シグナルをそれが表しているステージへ届け、リトライを負荷として扱い、価値の低い作業を意図的に切り捨てる。
このように構築されたクローラーには、もう一つ持つ価値のある特性がある。あるサイトが苦戦し始めると、クローラーはそれに気づき、自ら手を緩める。それはそのサイトにとって良いことであり、偶然ではなく、そのクローラーが来週もそこで歓迎され続けられる見込みにとっても良いことだ。