构建一个使用 Scrapeless 和 DuckDB 的爬取数据分析管道
Advanced Data Extraction Specialist
TL;DR:
- 一条以 JSON 文件为结束的抓取管道只是一半的管道;当数据落地后,分析问题才会出现。
- DuckDB 直接读取 JSON Lines,因此抓取的记录成为可以查询的表,没有服务器,没有架构迁移,也没有 ETL 工具。
- 这个管道通过 Scrapeless Universal Scraping API 获取 3 页,提取 30 条记录,将其加载到 DuckDB,并写入 Parquet。
- 使用 ZSTD 压缩的 Parquet 格式将同样的 30 条记录存储在 4,674 字节中,而 JSON Lines 则为 7,690 字节,并且 DuckDB 直接查询文件而不需要重新加载。
- 下面每个阶段都在你自己的机器上运行,无需云数据仓库账户和除 Scrapeless 密钥之外的任何凭证。
- 从 Scrapeless 免费计划 开始,并将抓取阶段指向你自己的来源。
管道概览
流程分为五个阶段,其中一个有趣的设计决策是数据停止作为文本并开始作为表的地方。
fetch (Scrapeless) → discover records → extract fields → transform to a typed table (DuckDB) → store as Parquet
JSON Lines 是抓取部分和分析部分之间的交接格式。它支持追加操作,能够在崩溃的运行中存活而不会损坏之前的数据,并且 DuckDB 本地读取它——因此不需要编写加载器。该格式在 JSON Lines 规范 中进行了说明。
这里没有任何内容需要数据仓库账户。DuckDB 在进程中运行,这使得它成为处理抓取数据的最便宜的方式。当目标是管理数据仓库时,Snowflake 数据摄取指南 介绍了相同的收集阶段,只是着陆区不同。
先决条件
- Python 3.9 或更高版本。
- 从仪表板获得的 Scrapeless API 密钥。
- 安装
duckdb:
bash
pip install "duckdb==1.5.4"
设置密钥:
bash
export SCRAPELESS_API_KEY="your_api_key_here"
第 1–3 阶段:抓取、发现、提取
通过 Scrapeless Universal Scraping API 每页一次调用返回渲染的 HTML,标准库解析器将每个引用块转换为一条记录。提取器每行发出一个 JSON 对象:
python
import json, os, urllib.request
from html.parser import HTMLParser
API = "https://api.scrapeless.com/api/v2/unlocker/request"
def fetch(url: str) -> str:
payload = json.dumps({
"actor": "unlocker.webunlocker",
"input": {"url": url, "js_render": True, "headless": True},
}).encode()
req = urllib.request.Request(
API, data=payload,
headers={"x-api-token": os.environ["SCRAPELESS_API_KEY"],
"Content-Type": "application/json"},
)
with urllib.request.urlopen(req, timeout=120) as r:
return json.loads(r.read())["data"]
class QuoteParser(HTMLParser):
def __init__(self):
super().__init__()
self.rows, self._cur, self._cap = [], None, None
def handle_starttag(self, tag, attrs):
a = dict(attrs); cls = a.get("class", "")
if tag == "div" and "quote" in cls:
self._cur = {"text": "", "author": "", "tags": []}
elif self._cur is not None and tag == "span" and "text" in cls:
self._cap = "text"
elif self._cur is not None and tag == "small" and "author" in cls:
self._cap = "author"
elif self._cur is not None and tag == "a" and "tag" in cls:
self._cap = "tag"
def handle_data(self, data):
if self._cap == "text": self._cur["text"] += data
elif self._cap == "author": self._cur["author"] += data
elif self._cap == "tag": self._cur["tags"].append(data.strip())
def handle_endtag(self, tag):
if self._cap in ("text", "author", "tag"): self._cap = None
if tag == "div" and self._cur and self._cur["text"]:
self.rows.append(self._cur); self._cur = None
rows = []
for page in range(1, 4):
p = QuoteParser(); p.feed(fetch(f"https://quotes.toscrape.com/page/{page}/"))
rows.extend({"page": page, **r} for r in p.rows)
print(f"抓取页面: 3 | 提取记录: {len(rows)}")
with open("quotes.jsonl", "w", encoding="utf-8") as f:
for r in rows: f.write(json.dumps(r, ensure_ascii=False) + "\n")
print(f"写入 quotes.jsonl ({os.path.getsize('quotes.jsonl')} 字节)")
text
抓取页面: 3 | 提取记录: 30
写入 quotes.jsonl (7690 字节)
这里的两个选择值得保留在你自己的版本中。
页面编号在提取时附加到每个记录上。出处的记录几乎是免费的,而后期重构的成本很高,它使你能够回答“这来自哪一页”而无需重新运行任何操作。
ensure_ascii=False 保持排版引号字符的完整性,而不是将其转义。当文本就是数据时,这一点很重要——转义输出依然是有效的 JSON,但它膨胀了文件大小,使后续检查变得更加困难。
第4阶段:转换为类型表
DuckDB 直接读取 JSON Lines 文件。read_json_auto 推断模式,周围的 SELECT 是你施加你实际想要的类型和派生列的地方:
python
import duckdb
con = duckdb.connect("quotes.duckdb")
con.execute("""
CREATE OR REPLACE TABLE quotes AS
SELECT
page::INTEGER AS page,
text AS quote_text,
author,
tags,
len(tags) AS tag_count
FROM read_json_auto('quotes.jsonl')
""")
total = con.sql("SELECT count(*) FROM quotes").fetchone()[0]
authors = con.sql("SELECT count(DISTINCT author) FROM quotes").fetchone()[0]
print(f"加载的行数: {total} | 不同作者数: {authors}")
print(con.sql("""
SELECT author, count(*) AS quotes, round(avg(tag_count), 2) AS avg_tags
FROM quotes
GROUP BY author
ORDER BY quotes DESC, author
LIMIT 5
""").to_df().to_string(index=False))
con.execute("COPY quotes TO 'quotes.parquet' (FORMAT PARQUET, COMPRESSION ZSTD)")
import os
print(f"parquet 字节数: {os.path.getsize('quotes.parquet')} | jsonl 字节数: {os.path.getsize('quotes.jsonl')}")
rt = duckdb.sql("SELECT count(*) AS n, count(DISTINCT author) AS a FROM 'quotes.parquet'").fetchone()
print(f"从 parquet 进行回传 -> 行数: {rt[0]}, 作者: {rt[1]}")
text
加载的行数: 30 | 不同作者数: 20
作者 引用数 平均标签数
阿尔伯特·爱因斯坦 6 2.83
J.K. 罗琳 3 1.33
鲍勃·马利 2 1.00
西斯博士 2 2.00
玛丽莲·梦露 2 4.00
parquet 字节数: 4674 | jsonl 字节数: 7690
从 parquet 进行回传 -> 行数: 30, 作者: 20
tags 列保留为列表,而不是扁平化为分隔字符串。DuckDB 将嵌套类型携带到 Parquet,因此 len(tags) 在 SQL 中有效并且数组在往返中得以保留——在这个阶段扁平化为 "a,b,c" 是一种值得打破的有损习惯。
CREATE OR REPLACE TABLE 使加载具有幂等性。重新运行这一阶段从当前文件重建表,而不是附加重复数据,这正是你在迭代提取器时所希望的行为。
聚合是整个过程的重点:30 条记录,20 位不同作者,其中一位作者占了 6 条。这是一个针对表的单行查询,而不是针对 JSON 文件的烦琐循环。
准备好在你关心的源上运行这个吗?创建一个免费的 Scrapeless 账户 并在提取阶段交换 URL。
第5阶段:存储为 Parquet
相同的 30 条记录在 ZSTD 压缩的 Parquet 中占用 4,674 字节,而在 JSON Lines 中占用 7,690 字节。在这个样本中节省是适度的;关注的原因在于这种格式在大规模使用时的作用以及后续所能实现的目标。
Parquet 是列存储,存储每一列的值及其各自的编码,编码在 Apache Parquet 文件格式规范 中定义。查询涉及五个列中的两个时,仅读取那两个。这里使用的压缩编解码器在 Zstandard 压缩标准 中具体说明。
往返线是该阶段的自我测试。读取 Parquet 文件返回 30 行和 20 位不同作者,与其来源的表匹配——在任何写入文件让另一个系统读取的管道中,值得进行验证。
注意最终查询将 'quotes.parquet' 直接读取为一个表,没有导入步骤,也没有与数据库文件的打开连接。这是以 Parquet 结束爬虫管道的实际原因:输出可以被 DuckDB 和大多数其他分析引擎查询,正是在它存放的地方。
完整的管道
作为一个脚本运行,五个阶段短到可以在一个窗口中读取。这是可以复制的版本:
python
import json, os, urllib.request
from html.parser import HTMLParser
API = "https://api.scrapeless.com/api/v2/unlocker/request"
def fetch(url: str) -> str:
payload = json.dumps({
"actor": "unlocker.webunlocker",
"input": {"url": url, "js_render": True, "headless": True},
}).encode()
req = urllib.request.Request(
API, data=payload,
headers={"x-api-token": os.environ["SCRAPELESS_API_KEY"],
"Content-Type": "application/json"},
)
with urllib.request.urlopen(req, timeout=120) as r:
return json.loads(r.read())["data"]
class QuoteParser(HTMLParser):
def __init__(self):
super().__init__()
self.rows, self._cur, self._cap = [], None, None
def handle_starttag(self, tag, attrs):
a = dict(attrs); cls = a.get("class", "")
if tag == "div" and "quote" in cls:
self._cur = {"text": "", "author": "", "tags": []}
elif self._cur is not None and tag == "span" and "text" in cls:
self._cap = "text"
elif self._cur is not None and tag == "small" and "author" in cls:
self._cap = "author"
elif self._cur is not None and tag == "a" and "tag" in cls:
self._cap = "tag"
def handle_data(self, data):
if self._cap == "text": self._cur["text"] += data
elif self._cap == "author": self._cur["author"] += data
elif self._cap == "tag": self._cur["tags"].append(data.strip())
def handle_endtag(self, tag):
if self._cap in ("text", "author", "tag"): self._cap = None
if tag == "div" and self._cur and self._cur["text"]:
self.rows.append(self._cur); self._cur = None
rows = []
for page in range(1, 4):
p = QuoteParser(); p.feed(fetch(f"https://quotes.toscrape.com/page/{page}/"))
rows.extend({"page": page, **r} for r in p.rows)
with open("quotes.jsonl", "w", encoding="utf-8") as f:
for r in rows: f.write(json.dumps(r, ensure_ascii=False) + "\n")
print(f"已获取页面:3 | 已提取记录:{len(rows)}")
import duckdb
con = duckdb.connect("quotes.duckdb")
con.execute("""
CREATE OR REPLACE TABLE quotes AS
SELECT page::INTEGER AS page, text AS quote_text, author, tags, len(tags) AS tag_count
FROM read_json_auto('quotes.jsonl')
""")
rows_loaded = con.sql("SELECT count(*) FROM quotes").fetchone()[0]
authors = con.sql("SELECT count(DISTINCT author) FROM quotes").fetchone()[0]
print(f"已加载行:{rows_loaded} | 不重复的作者:{authors}")
con.execute("COPY quotes TO 'quotes.parquet' (FORMAT PARQUET, COMPRESSION ZSTD)")
print(f"Parquet 字节数:{os.path.getsize('quotes.parquet')} | JSONL 字节数:{os.path.getsize('quotes.jsonl')}")
rt = duckdb.sql("SELECT count(*), count(DISTINCT author) FROM 'quotes.parquet'").fetchone()
print(f"从 Parquet 轮转到行数:{rt[0]},作者数:{rt[1]}")
## 接下来管道将去哪里
**按采集日期进行分区** 一旦您反复运行它。写入 `data/dt=<collection-date>/quotes.parquet` 结构可以让查询跳过整个目录,而不是扫描历史记录。
**保留 JSON Lines 文件。** 它们是所收集内容的原始记录;Parquet 是派生的、带类型的工件。当提取器错误出现时,您可以从原始文件重新推导,而不是重新抓取。
**在各个阶段之间添加断言。** 在提取和加载之间进行计数检查,可以捕捉到一个选择器悄然停止匹配的情况 — 否则会在几周后表现为行数悄然下降的故障模式。
在指向实时数据源之前,检查其条款和 `/robots.txt` 指令,这符合 <a href="https://datatracker.ietf.org/doc/html/rfc9309" rel="nofollow"><strong>机器人排除协议标准</strong></a>。限制收集到公共页面和像上述那样的边界页面范围。
## 结论
抓取器与分析管道之间的差距比看起来要小。JSON Lines作为交接,DuckDB作为查询引擎,以及Parquet作为存储工件,仅需一个依赖项且没有基础设施 — 3个页面进入,30行针对类型输出,可以直接查询。
大多数管道跳过的最后一步。写入 Parquet 然后读取回去以确认计数匹配,将存储步骤变成验证步骤,这就是您拥有的文件和可以信任的文件之间的区别。
[从 Scrapeless 免费计划开始](https://app.scrapeless.com/passport/login?utm_source=website&utm_medium=blog&utm_campaign=universalscrapingapi&utm_term=scrapeless-duckdb-parquet-pipeline) 对您的目标运行抓取阶段,并在对定期作业进行大小评估时审查 [Scrapeless 定价](https://www.scrapeless.com/zh/pricing?utm_source=website&utm_medium=blog&utm_campaign=universalscrapingapi&utm_term=scrapeless-duckdb-parquet-pipeline)。
## 常见问题
**问:为什么选择 DuckDB 而不是云仓库来存储抓取的数据?**
DuckDB 在进程内运行,无需服务器、账户和网络往返,这正符合大多数爬虫项目实际操作的规模。它本地支持 JSON 和 Parquet,因此无需编写加载器。当多个团队需要同时访问同一表时,云仓库才有其存在的价值——而不是当某个管道需要对其输出执行 SQL 时。
**问:在加载之前,我需要将嵌套字段(如标签列表)展平吗?**
不需要,展平会丢失信息。DuckDB 支持端到端的列表类型,因此 `tags` 在加载过程中仍然是一个数组,在 SQL 中通过 `len(tags)` 处理,并在 Parquet 的读取和写入中保持不变。将其压缩成分隔字符串会强迫之后的每个查询重新解析它。
**问:为什么在爬取和加载之间写入 JSON Lines?**
它实现了两个部分的解耦。爬取是一个缓慢且容易出错的部分;一旦记录以每行为线的格式存储在磁盘上,您可以在不重新获取的情况下,多次重新运行加载和转换。逐行追加也意味着,如果运行过程中断,之前的记录仍然完整且可读。
**问:Parquet 相对于 JSON Lines 小多少?**
在这个 30 条记录的样本中,4,674 字节对比 7,690 字节——大约小 40%。在这个大小下,不要过度解读这个比例,因为文件开销主导了结果。Parquet 的真正优势是列式读取:一个查询如果触及五列中的两列,只会读取那两列,这在文件不再能舒适地放入内存时尤为重要。
**问:我可以在不先将 Parquet 文件加载到数据库中的情况下查询它吗?**
可以,加载阶段的最后一行正是这样做的——`SELECT ... FROM 'quotes.parquet'`,没有打开数据库连接也没有导入。这使得 Parquet 成为爬虫管道的一个良好最终产物:输出保持可查询状态,DuckDB 和其他分析引擎均可访问。
在Scrapeless,我们仅访问公开可用的数据,并严格遵循适用的法律、法规和网站隐私政策。本博客中的内容仅供演示之用,不涉及任何非法或侵权活动。我们对使用本博客或第三方链接中的信息不做任何保证,并免除所有责任。在进行任何抓取活动之前,请咨询您的法律顾问,并审查目标网站的服务条款或获取必要的许可。



