最悪のスクレイパー障害は、障害のようには見えない。ジョブは実行され、リクエストは成功し、行は予定通りウェアハウスに到着する。その後、経理担当者が平均価格が一晩で3倍になった理由や、カタログの半分が突然在庫切れになった理由を尋ねてきて、調査の結果、3週間前に誰も気づかなかったウェブサイトのリデザインに行き着く。
これがスキーマドリフトだ。データは流れ続けるが、その形や意味が静かに変化している。これはウェブデータの通常の障害モードである。なぜなら、ソースは予告なしに変化するのに、エクストラクターは出力を生成し続けるからだ。本ガイドでは、これを検知する2つの層、不正な値を拒否するレコードレベルの契約と、データ全体がそれ自身らしさを失ったことに気づくバッチレベルのプロファイル、を扱う。また、小さなテストを通じて、なぜ両方が必要なのかを示す。
主なポイント
- ドリフトは通常、静かに起こる。変更されたページがエクストラクターを壊すことはめったにない。代わりに、エクストラクターにもっともらしく間違ったものを返させる。
- レコードレベルの契約は、各レコードを必須フィールド、型、範囲、形式のルールと照合する。明らかな破損を捕捉する。
- バッチレベルのプロファイルは、各実行をベースラインと比較する。充填率、型の構成、中央値などだ。契約では見えない破損を捕捉する。
- 3つの現実的なドリフトをテストしたところ、契約は1つを捕捉した。残りの2つはすべてのレコードレベルチェックを通過し、バッチプロファイルのみが検知した。
- データが出荷される前にドリフトに対してアラートを出し、バッチを隔離し、どのサイトとエクストラクターのバージョンがそれを生成したかを記録すること。
ドリフトが実際に起こる仕組み
| 原因 | データに起こること |
|---|---|
| リデザインによるマークアップの変更 | セレクタが別の要素にマッチする、または何にもマッチしない |
| 価格形式の変更 | ”9.99” が “€9.99”、“9,99”、あるいは最小単位での999になる |
| フィールドがJavaScriptの背後に移動 | 静的HTMLにはもう含まれず、空で返ってくる |
| ローカライゼーションまたは地域バリアント | 一部のページで異なる通貨、言語、単位が現れる |
| A/Bテスト | 一部のページが新しいレイアウトを使い、一部のレコードが壊れる |
| ブロックまたはチャレンジページ | ページは読み込まれ、エクストラクターは実行されるが、何も返さないかゴミを返す |
共通しているのは、これらのいずれも例外を発生させないことだ。最後の項目、正常に読み込まれるが求めていたページではないケースについては、the silent failure rateで扱っている。残りは、エクストラクターが自分の下で変化したページを忠実に処理しているだけだ。
層1: レコードレベルの契約
データ契約は、有効なレコードがどのようなものかを定義する。どのフィールドが必須か、その型、許容される範囲と形式だ。すべてのレコードは受け入れられる前にチェックされ、違反は静かに保存されるのではなく、カウントされ隔離される。
import re
import statistics
from collections import Counter
CONTRACT = {
"name": {"type": str, "required": True},
"price": {"type": float, "required": True, "min": 0.01, "max": 100_000},
"currency": {"type": str, "required": True, "pattern": r"^[A-Z]{3}$"},
"in_stock": {"type": bool, "required": False},
}
def violations(record, contract=CONTRACT):
"""Record-level checks: presence, type, range, format."""
problems = []
for field, rule in contract.items():
value = record.get(field)
if value in (None, ""):
if rule.get("required"):
problems.append(f"{field}: missing")
continue
if rule["type"] is float and isinstance(value, int) and not isinstance(value, bool):
value = float(value)
if not isinstance(value, rule["type"]):
problems.append(f"{field}: expected {rule['type'].__name__}, got {type(value).__name__}")
continue
if "min" in rule and value < rule["min"] or "max" in rule and value > rule["max"]:
problems.append(f"{field}: {value} out of range")
if "pattern" in rule and not re.match(rule["pattern"], value):
problems.append(f"{field}: {value!r} bad format")
return problems
{"name": "x", "price": "€9.99", "currency": "eur"} のようなレコードは2回失敗する。priceが文字列であることと、currencyが3文字の大文字コードでないことだ。契約は安価で、明示的で、理解しやすい。その限界は、各レコードが単独で判断されることであり、多くのドリフトは個々に見れば有効なレコードを生み出す。
層2: バッチレベルのプロファイル
プロファイルはバッチ全体を要約する。各フィールドについて、どのくらいの頻度で埋まっているか、どの型が現れるか、そして数値については中央値だ。各実行のプロファイルを最近の健全な実行から得たベースラインと比較することで、たとえすべてのレコードが契約に合格していても、データ全体が形を変えたときにそれがわかる。
def profile(records, fields=CONTRACT):
"""Batch-level shape: how often each field is filled, its types, and numeric medians."""
n = max(1, len(records))
out = {}
for field in fields:
values = [r.get(field) for r in records]
present = [v for v in values if v not in (None, "")]
numbers = [float(v) for v in present if isinstance(v, (int, float)) and not isinstance(v, bool)]
out[field] = {
"fill_rate": len(present) / n,
"types": Counter(type(v).__name__ for v in present),
"median": statistics.median(numbers) if numbers else None,
"distinct": len(set(map(str, present))),
}
return out
def drift(baseline, current, fill_drop=0.1, median_ratio=3.0):
"""Compare two batch profiles and describe what changed shape."""
alerts = []
for field, base in baseline.items():
cur = current[field]
if base["fill_rate"] - cur["fill_rate"] > fill_drop:
alerts.append(f"{field}: filled {base['fill_rate']:.0%} -> {cur['fill_rate']:.0%}")
if set(cur["types"]) - set(base["types"]):
alerts.append(f"{field}: new types {sorted(set(cur['types']) - set(base['types']))}")
if base["median"] and cur["median"]:
ratio = cur["median"] / base["median"]
if ratio > median_ratio or ratio < 1 / median_ratio:
alerts.append(f"{field}: median {base['median']:g} -> {cur['median']:g}")
return alerts
これらのしきい値は出発点にすぎない。充填率の10ポイントの低下や中央値の3倍の変化は、業務上の通常事態であることはめったにない。数週間分の履歴が溜まったら、フィールドごとに両方を調整すること。
なぜ両方が必要なのか: 小さなテスト
1,000件の商品レコードからなる健全なベースラインを生成し、3つの現実的な破損をシミュレートして、両方の層をそれぞれに対して実行した。
| ドリフト | 何が起こったか | レコードレベルの契約 | バッチプロファイル |
|---|---|---|---|
| フォーマットされた価格 | リデザインにより、ページの30%で通貨記号が価格の中に入り込んだ | 捕捉: 300件のレコードが拒否された | 捕捉: priceに新しい型strが出現 |
| 最小単位 | 価格が63.05ではなく6305として届き始めた | 見逃し: すべてのレコードが通過した | 捕捉: 中央値が63.05から6305へ、新しい型intも出現 |
| 壊れた在庫セレクタ | in-stockフィールドがページの80%で空になって返ってきた | 見逃し: そのフィールドは任意項目 | 捕捉: 充填率が100%から20%へ |
契約は、明らかに不正な形式の値を生み出した1つのドリフトを捕捉した。他の2つは、範囲内の正の数値や空のままの任意フィールドといった、個別には有効なレコードを生み出し、バッチビューだけが何かが変化したことを示した。最小単位のケースは仮定の話ではない。先週調査した実際の商品エンドポイントは、100.00ドルの靴に対して10000を返していた。詳しくはstop parsing HTMLを参照。
ドリフトが発火したときにすべきこと
- バッチを隔離する。 誰かが確認するまで、ドリフトアラートが出たバッチからのデータを出荷しないこと。遅れたデータセットの方が、間違ったデータセットよりましだ。
- 範囲を特定する。 アラートをサイト、ページタイプ、エクストラクターのバージョンごとに分解する。ドリフトは通常、1つの変更の後に1つのサイトで始まる。
- サンプルを比較する。 影響を受けたレコードのいくつかを、それらが由来するページと並べて見る。原因は通常、数分以内に明らかになる。
- 修正してバックフィルする。 エクストラクターを更新し、保存済みのページがあれば影響を受けた期間を再抽出し、なければ再収集する。
- 意図的にベースラインを更新する。 変更が正当なものであるとき、たとえばサイトが本当に通貨を変更した場合は、意図的にベースラインをリセットする。決して自動的には行わないこと。
そもそもドリフトを起こりにくくする
一部のソースは他のソースよりドリフトが少ない。検索エンジン向けに埋め込まれた構造化データは、ページレイアウトよりもはるかに頻度低く変化する。これが、それを先に抽出することで契約が発火する頻度全体を減らせる理由だ。ページ自体を監視することも役立つ。change detection at scaleによって検知されたテンプレートの変更は、抽出がまもなく壊れる可能性があるという早期警告であり、1つのサイトで契約違反率が上昇することは、a target health scoreへの強力な入力となる。1つのジョブにつき1つの市場から一貫して収集することも、地域や通貨のバリアントによって引き起こされるドリフトのクラス全体を取り除く。
結論
ウェブデータがドリフトするのは、ウェブが承諾なしに変化するからだ。危険なのはエクストラクターが壊れることではなく、それらがもはや構築時に想定されていた姿ではなくなったページ上で動き続け、1レコードずつ見れば問題なさそうなデータを生成し続けることだ。
すべてのレコードを契約と照合し、すべてのバッチをそれ自身の最近の履歴と照合すること。私たちのテストでは、契約だけでは3つのドリフトのうち1つしか捕捉できなかったが、バッチプロファイルは3つすべてを捕捉した。両方を組み合わせることで、「ダッシュボードが3週間おかしく見えていた」という事態を、それが始まった朝のアラートに変えることができる。
出典と参考文献
- 2026年9月29日にShifterが上記のコードを用いて実行したテスト。1,000件の生成された商品レコードと3つのシミュレートされたドリフトを使用。