大型采集系统架构与工程化
脚本能跑通不等于系统能长期运行。分界线在于:任务失败后能不能自动恢复、站点改版后能不能快速定位、数据能不能追溯来源、新增一个源要改多少代码。本章给出演进路径、模块划分、去重与分层设计、失败补偿、可测试性,最后附一张上线前检查清单。
演进路径:别一步跨到分布式
| 阶段 | 形态 | 规模参考 | 主要矛盾 | 触发升级的信号 |
|---|---|---|---|---|
| 一、单机脚本 | 一个 .py 跑到底 | 单源、每天几千条 | 快速验证可行性 | 需要多源、需要定时 |
| 二、任务表 + 多 worker | MySQL/Redis 存任务,多进程消费 | 数源、每天十万条 | 失败恢复与并发控制 | 单机带宽或 CPU 成瓶颈 |
| 三、分布式 + 代理 | scrapy-redis 多机、代理池、全局限速 | 多源、每天百万条 | 调度一致性、被封风险 | 需要对内提供稳定数据服务 |
| 四、平台化 | 配置驱动的平台 + 监控告警 + 数据分层 | 数十源、持续运行 | 工程规范与合规审计 | — |
绝大多数项目到第二阶段就够用。在没有任务表和失败恢复之前就上分布式,只会把「偶发丢数据」变成「随机丢数据」,更难查。
模块划分
| 模块 | 职责 | 关键接口 | 不该做的事 |
|---|---|---|---|
| fetcher | 发请求、限速、重试、代理 | fetch(url) -> Response | 不解析、不去重 |
| parser | HTML/JSON → 结构化字段 | parse(html) -> Item | 不发请求、不写库 |
| pipeline | 清洗、校验、去重、落地 | save(items) -> int | 不做网络请求 |
| scheduler | 触发、锁、并发控制 | run(task) | 不关心数据字段 |
| monitor | 指标、日志、告警 | report(metrics) | 不影响主流程 |
收益是可测试:parser 能拿固定 HTML 单测,pipeline 能拿构造数据单测,只有 fetcher 与 scheduler 依赖真实环境。把业务规则全塞在 fetch 与 parse 混杂的脚本里,是后期难维护的根因。
去重与增量设计
| 手段 | 计算方式 | 解决的问题 | 代价 |
|---|---|---|---|
| URL 指纹 | md5(url) | 同一页面被采多次 | 极低 |
| 规范化 URL 指纹 | 去掉 utm_*、?from= 后再取 md5 | 同一页面不同追踪参数 | 低,需维护参数黑名单 |
| 内容指纹 | md5(正文规范化后) | 同内容不同 URL(转载) | 低 |
| 相似指纹 | SimHash / MinHash | 洗稿、轻微改动 | 中,需阈值调参 |
| 更新时间戳 | 记录源的最新发布时间 | 增量抓取,减少请求量 | 低,源时间乱序时需缓冲窗口 |
| 业务主键 | sku_id、job_id | 跨源跨天幂等 | 需确认源侧主键稳定性 |
# fingerprint.py —— Python 3.10+
import hashlib, re
from urllib.parse import urlencode, urlparse, urlunparse
TRACKING = {"utm_source", "utm_medium", "utm_campaign", "from", "spm", "ref", "share_token"}
def normalize_url(url: str) -> str:
"""去掉追踪参数与锚点,保证同一页面的指纹稳定"""
parsed = urlparse(url)
query = [(k, v) for k, v in
(item.split("=", 1) if "=" in item else (item, "")
for item in parsed.query.split("&") if item)
if k not in TRACKING]
return urlunparse(parsed._replace(query=urlencode(sorted(query)), fragment=""))
def url_fingerprint(url: str) -> str:
return hashlib.md5(normalize_url(url).encode("utf-8")).hexdigest()
def content_fingerprint(text: str) -> str:
"""正文指纹:去空白与标点后再哈希,降低排版差异带来的误判"""
return hashlib.md5(re.sub(r"[\s\W_]+", "", text or "").encode("utf-8")).hexdigest()
if __name__ == "__main__":
a = "https://example.com/news/1.html?utm_source=wechat&from=timeline"
print(url_fingerprint(a) == url_fingerprint("https://example.com/news/1.html")) # True
数据分层与失败补偿
数据分三层:原始层 raw(响应体、状态码、抓取时间,保留 7~90 天)、清洗层 clean(规范化字段,可由 raw 重算)、应用层 app(汇总表、看板数据集)。分层的价值在于改解析规则不用重新请求站点:站点改版时拿原始层重跑 parser 即可,既省流量又不增加对方压力;原始层必须记录 url、status、fetched_at、raw_body,否则无法重算。
| 失败类型 | 处理方式 | 补偿机制 |
|---|---|---|
| 网络超时/5xx | 指数退避重试 2~3 次 | 仍失败写入失败表 |
| 解析失败 | 不重试(重试也不会变) | 落盘失败样本,修选择器后从 raw 层重跑 |
| 入库失败 | 整批回滚,保留该批数据 | 失败批次落 JSONL,确认后重放 |
| 任务中断 | 任务表状态停在「进行中」 | 按 update_time 超时回收,重新入队 |
| 数据缺口 | 事后发现某天条数异常 | 按日期区间补跑,靠业务主键幂等写入 |
补偿脚本的做法很简单:从失败文件里取前 N 条任务,逐条调用 handler 并限速,成功即移除、失败带错误原因写回文件。前提是 handler 幂等——靠第 18 章唯一键 + ON DUPLICATE KEY UPDATE、第 19 章 UpdateOne(upsert=True),重放不会产生重复数据;不幂等的系统只能「小心地再跑一遍」,等于把风险留给下一个人。
可测试性:用固定样例守住解析器
# test_parser.py —— 离线单测:pytest test_parser.py
from parser import parse_detail
SAMPLE_HTML = """
<html><body><h1>示例标题</h1><span class="time">2024-05-01 09:00</span>
<div class="article-content">
<p>第一段正文,长度足够用于通过阈值校验。</p>
<p>第二段正文,同样用于测试段落合并逻辑是否正常。</p>
<div class="recommend">推荐阅读:不应出现在正文里</div>
</div></body></html>
"""
def test_parse_detail_basic_fields():
item = parse_detail(SAMPLE_HTML, "https://example.com/1")
assert item is not None and item["title"] == "示例标题"
assert "第一段正文" in item["content"]
def test_recommend_block_is_removed():
assert "推荐阅读" not in parse_detail(SAMPLE_HTML, "https://example.com/1")["content"]
def test_short_content_returns_none():
"""正文过短应判为解析失败,而不是产出脏数据"""
html = "<html><h1>标题</h1><div class='article-content'><p>太短</p></div></html>"
assert parse_detail(html, "https://example.com/2") is None
站点改版后把新 HTML 存成样例,先跑测试看哪条断言挂了,比重新读代码快得多;每个源至少准备三份样例:正常页、边界页(无作者/无时间)、异常页(结构缺失)。
容量评估与「什么时候该停」
估算方式:请求量 ≈ 页面数 ÷ 每页链接数 × 分页数,带宽 = 平均页面大小 × 请求量,存储 = 单条大小 × 日增量 × 保留天数,解析成本 = 每条 CPU 时间 × 日增量。请求量超过目标站日均流量 1% 就不礼貌,该考虑官方接口了。
判断「该停」的三个信号:单源维护成本长期高于数据价值;站点已明确拒绝或采取技术措施阻止;用途无法说清合法性依据。采集系统最重要的工程决策之一,是明确哪些源不做。
常见坑
| 现象 | 原因 | 处理 |
|---|---|---|
| 站点改版后数据为空却没人发现 | 只监控进程存活,没监控产出 | 用「条数突降」做静默失败检测 |
| 补数据后出现重复 | 写入不幂等 | 唯一键 + upsert,重跑才安全 |
| 分布式后故障更难查 | 日志缺少统一标识 | 日志带 site/trace_id,指标按源打标签 |
| 新增一个源要改主流程 | 站点差异硬编码在代码里 | 选择器与限流参数移到配置 |
上线前检查清单
| 类别 | 检查项 | 通过标准 |
|---|---|---|
| 合规 | 数据来源与用途 | 只采公开数据,遵守 robots.txt 与站点条款,保留来源链接 |
| 合规 | 隐私与访问控制 | 不采集个人隐私信息,不绕过登录、验证码、付费墙 |
| 合规 | 使用边界 | 再分发前取得授权,符合《个人信息保护法》《数据安全法》《著作权法》 |
| 稳定性 | 限速与并发 | 有全局 QPS 上限与退避策略,频率可控 |
| 稳定性 | 幂等 | 唯一键 + upsert,重跑不产生重复 |
| 稳定性 | 失败恢复 | 有失败表与重跑脚本,任务中断可回收 |
| 可观测 | 指标与日志 | 成功率、队列、解析失败率、入库条数四项齐备 |
| 可观测 | 静默失败检测 | 条数突降有告警,且与日志口径一致 |
| 可维护 | 配置与测试 | 站点差异在配置,parser 有固定 HTML 样例单测 |
| 可追溯 | 数据分层 | raw 层可重算,clean 层字段有清洗记录 |
小结:架构演进按「单机脚本 → 任务表加多 worker → 分布式代理 → 平台化」逐级推进,每一级都由真实瓶颈触发;把系统拆成 fetcher、parser、pipeline、scheduler、monitor 五个模块,让业务规则集中且可单测;去重靠 URL 与内容指纹、增量靠更新时间戳,数据分原始层与清洗层以便改规则时重算;失败一律走「幂等 + 补偿」;最后一条同样重要——明确哪些源不做,把合规边界当成架构的一部分。