AirflowとScrapelessを使ってスケジュールされたウェブスクレイピングパイプラインを構築する方法
Lead Scraping Automation Engineer
TL;DR:
- Airflow におけるスクレイピングパイプラインは、取得、抽出、保存という3つのタスクで構成されており、戻り値を渡すことで接続されています。このフルランはここで 25.2 秒で完了しました
state=success。 - Airflow 3 ではスケジューリングパラメータの名前が変更されました。
schedule_interval=はTypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'を引き起こすため、Airflow 2 用に書かれたほとんどのチュートリアルは DAG 定義で止まります。 - 制約ファイルを使って Airflow をインストールします。これがないと、
pip install apache-airflow==3.1.2は依存関係の解決にここで 13 分以上かかりましたが、制約ファイルを使用すると完了しました。 - DAG を実行するのに Docker は必要ありません。
airflow dags testは SQLite に対してプロセス内でグラフ全体を実行します。 - タスク間で渡される値はメタデータデータベースにシリアライズされます — 取得タスクの HTML は、解析されたレコードの 2,154 バイトに対して 53,889 バイトとして保存されました。
- 同じ DAG を二回実行すると 40 行ではなく 20 行が残ります。なぜなら、書き込みは製品タイトルにキーを基づくアップサートだからです。
- Scrapeless 無料プラン のレンダリングされたページに対してスケジューリングを開始します。
実行することを思い出したときだけ走るスクレイパーはスクリプトです。毎朝日付の付いたテーブルを生成する何かに変えるということは、2回目の実行が昨日の行にどのような影響を与えるのか、資格情報がどこにあるのかを決定することを意味します。Airflow はスケジューリングとタスクの接続を処理します; これら2つの決定はあなたのものです。
以下のパイプラインは、公共図書カタログを日次スケジュールで収集します:ユニバーサルスクレイピングAPIを通じてレンダリングされた取得、レコードへのパース、SQLite へのアップサート。すべての数字は Airflow 3.1.2 での実行から来ています。
Pipeline at a Glance
| Stage | Task | Input | Output | Verified result |
|---|---|---|---|---|
| Fetch | fetch |
カテゴリー URL | レンダリングされた HTML | 50,403 文字 |
| Extract | extract |
HTML 文字列 | レコードのリスト | 20 レコード |
| Store | store |
レコードリスト | 行数 | SQLite に 20 行 |
フローは 取得 → 抽出 → 保存 であり、各ステージはその戻り値を次のステージに渡します。その3つのポイントで境界を保つことが、他のステージを繰り返すことなく1つのステージを再実行できることを可能にし、ネットワークコールなしでパースロジックをテスト可能に保ちます。
Install Airflow Without the Long Resolve
Airflow は大規模な固定グラフに依存しており、対応する制約ファイルなしでインストールすると、pip がバージョンの組み合わせをロールバックし続けます。ここでの単純な pip install apache-airflow==3.1.2 は、13 分後もまだ解決中でしたが、制約ファイルを使った同じインストールは通常通り完了しました。
bash
python3 -m venv .venv
./.venv/bin/pip install "apache-airflow==3.1.2" \
--constraint "https://raw.githubusercontent.com/apache/airflow/constraints-3.1.2/constraints-3.12.txt"
./.venv/bin/pip install "parsel==1.10.0" "requests==2.32.5"
制約の URL は Airflow のバージョンと Python のバージョンの両方をエンコードしており、プロジェクトのインストール参照は、最適化ではなくインストールの一部として扱います。AIRFLOW_HOME を明示的にどこかにポイントし、メタデータデータベースを作成します:
bash
export AIRFLOW_HOME="$PWD/home"
./.venv/bin/airflow db migrate
text
Airflow database tables created
DB: sqlite:////root/verify-airflow/home/airflow.db
Database migrating done!
SQLite はデフォルトのバックエンドであり、単一マシンのパイプラインには十分です。スケジューラはタスクを並行して実行するために同時データベースが必要ですが、以下のものはそれを必要としません。
Stage 1: Fetch a Rendered Page
最初のステージは、パイプラインの残りがデータを見るかどうかを決定します。単純な HTTP クライアントは、JavaScript が実行される前にサーバーが送信するものを返すため、取得タスクはユニバーサルスクレイピングAPIを呼び出し、それはレンダリングされた文書を返します。
python
@task
def fetch() -> str:
payload = json.dumps({
"actor": "unlocker.webunlocker",
"input": {"url": CATEGORY_URL, "js_render": True, "headless": False},
}).encode()
request = urllib.request.Request(
"https://api.scrapeless.com/api/v2/unlocker/request",
data=payload,
headers={"Content-Type": "application/json",
"x-api-token": os.environ["SCRAPELESS_API_KEY"]},
)
with urllib.request.urlopen(request, timeout=180) as response:
body = json.loads(response.read())
html = body["data"]
print(f"fetched {len(html)} chars")
return html
text
fetched 50403 chars
レスポンスエンヴェロープは {"code": ..., "data": ...} であり、data は HTML を文字列として保持しています。API が JSON を返すため、文書はすでに UTF-8 としてデコードされた状態で到着します — 各価格の £ はタスクの中でエンコーディング処理なしに生き残ります。
API トークンは DAG ファイルではなく、環境から来ます。Airflow はタスクを実行するときにプロセス環境を読み取るため、パイプラインを実行するシェルでエクスポートすると、リポジトリに資格情報が残りません。Airflow Variables および Connections は、複数の DAG がそれを必要とする場合のより完全な回答です。
Stage 2: Extract Records
抽出ステージは文字列を受け取り、レコードを返します。ネットワークに接触せず、マークアップが変更されるたびに保存されたフィクスチャに対して行えるため、テストが可能です。
python
@task
def extract(html: str) -> list[dict]:
sel = Selector(text=html)
rows = []
for card in sel.css("article.product_pod"):
rows.append({
"title": card.css("h3 a::attr(title)").get(),
"price": card.css("p.price_color::text").get(),
"rating": (card.css("p.star-rating::attr(class)").get() or "").replace("star-rating", "").strip(),
"in_stock": "In stock" in (card.css("p.instock.availability::text").getall() or [""])[-1],
})
print(f"extracted {len(rows)} records")
return rows
text
extracted 20 records
各カードを最初に選択し、その内部をクエリすることで、同じ製品のフィールドがまとめられます。すべてのタイトルとすべての価格を2つのフラットリストとして選択し、それらをジップすると、1つのカードがフィールドを欠く瞬間に静かに不一致のペアが生成されます。
Stage 3: Store With an Upsert
スケジュールされたパイプラインは、明日同じ行を書き込むため、書き込みは、盲目的に追加するのではなく、レコードを一意にするものとして定義する必要があります。
python
@task
def store(rows: list[dict]) -> int:
conn = sqlite3.connect(DB_PATH)
conn.execute("""CREATE TABLE IF NOT EXISTS books (
title TEXT PRIMARY KEY, price TEXT, rating TEXT,
in_stock INTEGER, scraped_at TEXT)""")
now = datetime.utcnow().isoformat(timespec="seconds")
conn.executemany(
"INSERT OR REPLACE INTO books VALUES (?,?,?,?,?)",
[(r["title"], r["price"], r["rating"], int(r["in_stock"]), now) for r in rows],
)
conn.commit()
count = conn.execute("SELECT COUNT(*) FROM books").fetchone()[0]
conn.close()
print(f"stored; table now holds {count} rows")
return count
text
stored; table now holds 20 rows
titleは主キーであり、SQLiteのINSERT OR REPLACEの動作は新しい行を挿入する前に競合する行を削除します。実行後の保存されたテーブル:
text
title price rating stock
A Murder in Time £16.64 One 1
A Study in Scarlet (Sherlock Holmes £16.73 Two 1
A Time of Torment (Charlie Parker #1 £48.35 Five 1
Boar Island (Anna Pigeon #19) £59.48 Three 1
Delivering the Truth (Quaker Midwife £20.89 Four 1
Hide Away (Eve Duncan #20) £11.84 One 1
ステージをDAGに接続する
DAGの宣言はスケジュールを持ち、最後の行は戻り値を渡すことで依存関係のチェーンを宣言します。
python
@dag(
dag_id="books_pipeline",
schedule="0 6 * * *",
start_date=pendulum.datetime(2026, 9, 1, tz="UTC"),
catchup=False,
tags=["scraping"],
)
def books_pipeline():
...
store(extract(fetch()))
books_pipeline()
scheduleは標準cron式を使用するため、0 6 * * *はDAGのタイムゾーンで毎日06:00です。start_dateを明示的なゾーンで設定することは、見た目以上に重要です:デイライトセービングを観察するゾーンに固定されたDAGは、年に2回実際の実行時間がずれるため、オフセットはIANAタイムゾーンデータベースから来ます。UTCに固定すると、間隔は固定されたままになります。catchup=Falseは特にスクレイパーにとって重要です:これがTrueに設定されていると、start_dateが1か月前のDAGを展開すると、逃した間隔ごとに実行をキューイングし、スクレイパーは3週間前のページを回復できません。
store(extract(fetch()))の作成は全体の依存関係グラフです。Airflowは呼び出しのネストを読み取り、エッジを構築します。
Dockerなしで実行する
ほとんどのAirflowのウォークスルーは、Docker ComposeファイルとPostgresコンテナから始まります。DAGを開発するために、airflow dags testはスケジューラー、ウェブサーバー、またはコンテナを起動せずに、プロセス内で全体のグラフを実行します。
bash
export AIRFLOW_HOME="$PWD/home"
export SCRAPELESS_API_KEY="your-api-key"
./.venv/bin/airflow dags test books_pipeline
text
fetched 50403 chars
extracted 20 records
stored; table now holds 20 rows
DagRun Finished: dag_id=books_pipeline, run_duration=25.211182, state=success
その実行は、パースロジックがまだ動いている間に利用可能な最速のフィードバックループです。DAGが期待通りに動いているときにスケジューラーが開始されます。
クライアント側でレンダリングされるページに対してスケジューリングしますか?Scrapelessの無料プランは、エンドツーエンドで毎日DAGを実行するために十分なリクエストをカバーしています。
2回実行することは行を倍にしないはずです
スケジュールされたパイプラインは無人で実行されるため、興味深いテストは最初の実行ではなく2回目の実行です。同じDAGを再度実行すると:
text
stored; table now holds 20 rows
20、ではなく40行ではありません。titleに基づくアップサートは各行を置き換え、そのscraped_atを更新し、これがデバッグ中に手動でDAGをトリガーするのに安全な理由です。単純なINSERTであれば、テーブルが倍になり、どのコピーが現在のものであるかを判断する方法がなくなってしまいます。
履歴が重要な場合—価格の動きを追跡する—キーはタイトルと実行日の日付のペアになり、タイトルだけではなく、昨日の行が残ります。
XCom境界
戻り値はXComを介してタスク間で移動し、XComの値はメタデータデータベースにシリアライズされます。ここで保存された3つのステージ:
text
task_id serialized bytes
fetch 53889
extract 2154
store 2
HTMLは解析されたレコードの25倍のサイズであり、それはメモリではなくAirflowデータベースにあります。このサイズでは問題ありませんが、ページが大きくなるか、実行が頻繁になると問題が発生します。
スケールする形は早期に解析し、小さく渡すことです:フェッチタスクが生のドキュメントをオブジェクトストレージに書き込み、キーを返し、次にエクストラクトがそれを読み取るのです。ステージの境界は同じままで、メタデータデータベースはドキュメントの代わりにレコードサイズの値を保持し続けます。
Airflow 2のコードはここで変更なしでは実行できません
2つの変更が古いチュートリアルに従っているすべての人を捕まえ、1つだけが自らを告知します。
python
from airflow.decorators import dag, task # deprecated
from airflow.sdk import dag, task # Airflow 3
旧いインポートパスはまだ解決され、DeprecatedImportWarning: The airflow.decorators.dag attribute is deprecated. Please use 'airflow.sdk.dag'を出力します。スケジューリングパラメーターはハードフェイルです:
text
TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'
schedule_interval=はschedule=になりました。ほぼすべての公開されたAirflowスクレイピングチュートリアルはバージョン3以前のものであるため、TypeErrorは多くの読者が最初に遭遇するものであり、これはDAGパース時間に現れ、それ以降のスクレイピングコードが実行される前に表示されます。
結論
このパイプラインは3つの関数と1つのデコレーターで構成されています:フェッチはHTMLを返し、エクストラクトはレコードを返し、ストアはカウントを返し、store(extract(fetch()))がグラフです。スクリプトではなくパイプラインとして機能する理由は、2回目の実行を生き延びるアップサート、cronとして表現されたスケジュール、再フェッチせずに再解析できるステージの境界です。
airflow dags test と SQLite から始め、文書が XCom から大きくなるときは除外し、インストールを制約ファイルで固定して、最初に戦うのは依存関係解決者ではなく、マークアップになります。同じ問題のストレージエンドについては、私たちの DuckDB と Parquet パイプライン が列指向出力をカバーし、料金 では、日次スケジュールのリクエストあたりのコストが示されています。
スケジュールでスクレイパーを設置する準備はできましたか? Scrapeless の無料プランから始めて フェッチステージを自分のターゲットに向けてください。
FAQ
Q: スクレイピングのために Airflow を実行するのに Docker は必要ですか?
いいえ。 airflow dags test <dag_id> はデフォルトの SQLite データベースに対してプロセス内でフルグラフを実行し、上記の 25.2 秒の実行がどのように行われたかを示しています。Docker と Postgres は、スケジューラーを無人で実行し、タスクを並行して実行したい場合に役立ちます。SQLite は必要な同時書き込みをサポートしていないためです。
Q: Airflow DAG の API キーはどこに置くべきですか?
DAG ファイルには入れないでください。プロセス環境から読み取ることで、バージョン管理から外し、いくつかの DAG が同じ資格情報を必要とする場合には、Airflow 変数または接続がより良い場所です。どちらもメタデータデータベースに保存され、コードに貼り付けられるのではなく名前で参照されます。
Q: スケジュールされたスクレイパーを重複行を作成することから止めるにはどうすればいいですか?
レコードをユニークにする要素を定義し、そのキーに対して書き込みます。ここでは title がプライマリキーで、書き込みが INSERT OR REPLACE であるため、2 回目の実行で 40 行ではなく 20 行が残りました。現在の状態ではなく履歴を望む場合は、キーを広げて実行日を含めることで、各日のスナップショットが自分自身の行になります。
Q: Airflow タスク間でスクレイピングした HTML を渡すことは安全ですか?
小さいサイズでは、はい — 上記の 50,403 文字のドキュメントはメタデータデータベースで 53,889 バイトにシリアル化され、タスク間で問題なく移動しました。ドキュメントや実行頻度が増えると、それは良いアイデアではなくなるため、すべての値がそのデータベースに存在します。ドキュメントをオブジェクトストレージに書き込み、キーを渡すことで、レコードサイズの XCom 値を持つ同じタスク境界を維持します。
Q: なぜ私の Airflow チュートリアルコードは schedule_interval で TypeError を発生させるのですか?
それは Airflow 2 用に書かれたためです。バージョン 3 ではパラメータの名前を schedule に変更し、古い名前を渡すと DAG が解析されている間に TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval' が発生します。伴う変更はインポートです: airflow.decorators はまだ機能しますが警告を出し、airflow.sdk が現在のパスです。
Q: スクレイピング用の DAG に catchup を有効にすべきですか?
通常はしません。 catchup=True では、start_date が過去にある DAG は、デプロイされるとすぐに欠落した間隔ごとに 1 回の実行をキューに入れます。その動作は日付が古いデータセットの再処理に適していますが、スクレイパーは現在ページに表示されているものを読み取るため、これらの実行は今日のデータを収集し、それに歴史的な日付を付与します。
Scrapelessでは、適用される法律、規制、およびWebサイトのプライバシーポリシーを厳密に遵守しながら、公開されているデータのみにアクセスします。 このブログのコンテンツは、デモンストレーションのみを目的としており、違法または侵害の活動は含まれません。 このブログまたはサードパーティのリンクからの情報の使用に対するすべての責任を保証せず、放棄します。 スクレイピング活動に従事する前に、法律顧問に相談し、ターゲットウェブサイトの利用規約を確認するか、必要な許可を取得してください。



