大型采集系统架构与工程化

采集系统架构:调度、队列、Worker 与数据分层

脚本能跑通不等于系统能长期运行。分界线在于:任务失败后能不能自动恢复、站点改版后能不能快速定位、数据能不能追溯来源、新增一个源要改多少代码。本章给出演进路径、模块划分、去重与分层设计、失败补偿、可测试性,最后附一张上线前检查清单。

演进路径:别一步跨到分布式

阶段形态规模参考主要矛盾触发升级的信号
一、单机脚本一个 .py 跑到底单源、每天几千条快速验证可行性需要多源、需要定时
二、任务表 + 多 workerMySQL/Redis 存任务,多进程消费数源、每天十万条失败恢复与并发控制单机带宽或 CPU 成瓶颈
三、分布式 + 代理scrapy-redis 多机、代理池、全局限速多源、每天百万条调度一致性、被封风险需要对内提供稳定数据服务
四、平台化配置驱动的平台 + 监控告警 + 数据分层数十源、持续运行工程规范与合规审计

绝大多数项目到第二阶段就够用。在没有任务表和失败恢复之前就上分布式,只会把「偶发丢数据」变成「随机丢数据」,更难查。

模块划分

模块职责关键接口不该做的事
fetcher发请求、限速、重试、代理fetch(url) -> Response不解析、不去重
parserHTML/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_idjob_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 即可,既省流量又不增加对方压力;原始层必须记录 urlstatusfetched_atraw_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 与内容指纹、增量靠更新时间戳,数据分原始层与清洗层以便改规则时重算;失败一律走「幂等 + 补偿」;最后一条同样重要——明确哪些源不做,把合规边界当成架构的一部分。

笔记加载中…