如何使用 Airflow 和 Scrapeless 构建定时网页抓取流水线
Lead Scraping Automation Engineer
TL;DR:
- 在 Airflow 中,抓取管道包含三个任务——获取、提取、存储——通过传递返回值连接,并且完整运行在这里耗时 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 在那里占据了 53,889 字节,而解析的记录只占 2,154 字节。
- 两次运行相同的 DAG 留下了 20 行而非 40 行,因为写入是基于产品标题的 upsert。
- 开始在 Scrapeless 免费计划 上调度渲染页面。
一个在你记得运行时才运行的抓取器是一个脚本。将其变成每个早晨生成带日期的表格意味着要决定第二次运行对昨日行的影响以及凭证存放的位置。Airflow 处理调度和任务连接;这两个决定仍然属于你。
下面的管道每天收集一个公共书籍目录:通过通用抓取 API 渲染的抓取、解析为记录、以及插入到 SQLite。每个数字来自于在 Airflow 3.1.2 上运行它的结果。
Pipeline at a Glance
| 阶段 | 任务 | 输入 | 输出 | 验证结果 |
|---|---|---|---|---|
| 获取 | fetch |
分类 URL | 渲染的 HTML | 50,403 字符 |
| 提取 | extract |
HTML 字符串 | 记录列表 | 20 条记录 |
| 存储 | store |
记录列表 | 行数 | SQLite 中有 20 行 |
流程是 获取 → 提取 → 存储,每个阶段将其返回值传递给下一个。保持在这三个点之间的界限使你可以重新运行一个阶段而不重复其他阶段,并使解析逻辑在没有网络调用的情况下可测试。
安装 Airflow 而无需长时间解析
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 是默认后端,足以满足单机管道的需求。调度程序需要一个并发数据库以并行运行任务,但下面的内容不需要这一点。
阶段 1:获取渲染页面
第一个阶段决定了管道的其余部分是否能够看到数据。一个普通的 HTTP 客户端在 JavaScript 执行之前返回服务器发送的任何内容,因此抓取任务调用 通用抓取 API,该 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 token 来自于环境而不是 DAG 文件。Airflow 在执行任务时读取进程环境,因此在运行管道的 shell 中导出它可将凭证保留在代码库之外。Airflow 变量和连接是更多 DAG 需要时的完整答案。
阶段 2:提取记录
提取阶段接收一个字符串并返回记录。它不涉及网络,这意味着无论何时标记发生变化,都可以针对保存的模板进行测试。
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
首先选择每个卡片,然后在其中查询可以将来自同一产品的字段放在一起。将所有标题和所有价格作为两个平面列表选取并压缩在一起会产生无声不匹配的对,只要一张卡片缺少某个字段。
阶段 3:使用 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 比看起来更重要:一个 pinned 到观察夏令时的时区的 DAG 每年两次调整其实际执行时间,而偏移量来自 IANA 时区数据库。将其固定在 UTC 中可以保持间隔不变。catchup=False 对于爬虫特别重要:当它设置为 True 时,部署一个 start_date 是一个月前的 DAG,会为每个错过的间隔排队运行,而爬虫无论如何无法恢复三周前的页面。
编写 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 从头到尾运行。
两次运行不应使您的行数翻倍
计划的管道是在无人值守的情况下运行,因此有趣的测试是第二次运行而不是第一次。再次执行相同的 DAG:
text
stored; table now holds 20 rows
二十,而不是四十。针对 title 的 upsert 替换了每一行并刷新了 scraped_at,这使得在调试时手动触发 DAG 是安全的。一个普通的 INSERT 将使表格翻倍,并且没有办法知道哪个副本是当前的。
如果历史很重要——跟踪价格如何变化——关键变成了标题和运行日期的组合,而不是单独的标题,昨天的行依然保留。
XCom 边界
返回值通过 XCom 在任务之间移动,XCom 值被序列化到元数据数据库中。这里存储的三个阶段:
text
task_id serialized bytes
fetch 53889
extract 2154
store 2
HTML 的大小是从中解析出的记录的 25 倍,并且它存储在 Airflow 数据库中而不是内存中。这个大小是可以接受的,但当页面变大或运行变得更频繁时,这就不再适用了。
可以扩展的形状是提前解析并传递小数据:让抓取任务将原始文档写入对象存储并返回一个密钥,然后让提取任务读取它。阶段边界保持不变,元数据数据库继续持有记录大小的值而不是文档。
Airflow 2 代码在这里不改变情况下无法运行
两个更改可能会让跟随旧教程的人感到困惑,只有其中一个会显示自己。
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 解析时浮出水面,在任何爬虫代码运行之前。
结论
该管道由三个函数和一个装饰器组成:fetch 返回 HTML,extract 返回记录,store 返回计数,而 store(extract(fetch())) 是图形。使其成为管道而不是脚本的是在第二次运行时存活的 upsert、以 cron 表达的调度,以及让您无需重新抓取即可重新解析的阶段边界。
从airflow dags test和SQLite开始,保持文档不溢出XCom,一旦它们变大,就通过约束文件固定安装,这样你第一件要解决的就是标记,而不是依赖解析器。对于同一问题的存储端,我们的DuckDB和Parquet管道涵盖了列输出,而定价列出了每日计划在请求中花费的成本。
准备好将抓取程序安排到计划上吗?从Scrapeless免费计划开始,并将提取阶段指向你自己的目标。
FAQ
问:我需要Docker来运行Airflow以进行抓取吗?
不需要。airflow dags test <dag_id>在默认的SQLite数据库中以进程执行完整图形,这就是上述25.2秒运行的执行方式。当你想要调度程序无监督地运行,并且任务并行执行时,Docker和Postgres就会变得有用,因为SQLite不支持需要的并发写入。
问:API密钥应该放在哪里在Airflow DAG中?
不应放在DAG文件中。从进程环境读取可以将其排除在版本控制之外,而Airflow变量或连接是多个DAG需要相同凭证时更好的存放位置——两者都存储在元数据数据库中,并通过名称引用,而不是粘贴到代码中。
问:如何停止调度抓取程序创建重复行?
定义使记录唯一的内容,并且根据该密钥写入。这里title是主键,写入是INSERT OR REPLACE,因此第二次运行留下了20行,而不是40行。如果你想记录历史而不是当前状态,可以扩大密钥以包含运行日期,这样每一天的快照都是自己的行。
问:在Airflow任务之间传递抓取的HTML是否安全?
在小尺寸时是安全的——上述50,403个字符的文档在元数据数据库中序列化为53,889字节,并在任务之间顺利移动。随着文档或运行频率的增长,这变得不再是个好主意,因为每个值都存在于该数据库中。将文档写入对象存储并传递一个密钥可以保持相同的任务边界,同时拥有记录大小的XCom值。
问:为什么我的Airflow教程代码在schedule_interval上引发TypeError?
因为它是为Airflow 2编写的。版本3将参数重命名为schedule,并在解析DAG时传递旧名称会引发TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'。伴随改变的是导入:airflow.decorators仍然有效但会发出警告,而airflow.sdk是当前路径。
问:抓取DAG是否应该启用catchup?
通常不需要。使用catchup=True,在过去的start_date的DAG在部署时会排队每个错过的时间间隔进行一次运行。这种行为适合于重新处理过时的数据集,但抓取程序读取的是页面现在显示的内容,因此这些运行将收集今天的数据并用历史日期进行标记。
在Scrapeless,我们仅访问公开可用的数据,并严格遵循适用的法律、法规和网站隐私政策。本博客中的内容仅供演示之用,不涉及任何非法或侵权活动。我们对使用本博客或第三方链接中的信息不做任何保证,并免除所有责任。在进行任何抓取活动之前,请咨询您的法律顾问,并审查目标网站的服务条款或获取必要的许可。



