长周期爬虫的数据一致性:断点续爬 + 事务回滚保障采集质量

简介: 长周期爬虫的数据一致性:断点续爬 + 事务回滚保障采集质量

一、重新定义问题:爬虫是一类有状态的流处理任务
长周期爬虫(任务跨度数天、采集千万级页面)在工程上不应被当成"脚本",而应被建模为一条有状态的流处理管道:待爬 URL 是输入流,落库记录是副作用(side effect),采集进度是消费位点(offset)。一旦这样抽象,文章标题背后的真问题就浮现了——
在「消费输入」与「产生副作用」两个动作之间,如何让进度推进与数据写入保持原子?
这正是 Kafka 消费者 exactly-once 语义的经典难题:若先提交 offset 再写库,崩溃会漏采;若先写库再提交 offset,崩溃会重采。长周期爬虫的"脏数据"全部源于这个原子性鸿沟。本文给出一套不依赖特定厂商、可在任意 RDBMS 上落地的解决模型。
二、从投递语义看本质
消息系统有三类投递语义,映射到爬虫:
at-most-once:写库失败不重试 → 漏采,不可接受。
at-least-once:重试到底 → 重采,是爬虫的宿命语义,因为网络天然不可靠。
exactly-once:理想态,现实中等价于 "at-least-once + 幂等收敛"。
结论先行:不要追求真正的 exactly-once(需分布式事务,代价过高),而是接受 at-least-once,用幂等键把重复副作用收敛为有效一次(effectively-once)。这是整套设计的第一性原理。
三、不一致的三个来源(精确表述)
现象
技术根因
归类
半条记录
单条记录跨多语句/多表写,进程在语句间被杀
非原子副作用
重复记录
offset 提交早于副作用落库,崩溃后同批重放
offset/side-effect 鸿沟
漏采 + 跳号
批次部分提交成功但进度按整批推进
检查点粒度粗于事务
三者统一根因:进度状态与业务数据不在同一事务边界内。
四、核心架构:单一事务边界(状态库即真相源)
关键设计决策——把"采集进度"也当作业务数据的一部分,与落库记录放在同一个本地 RDBMS 事务里提交。如此,外部消息队列只负责"投递 URL",不再承担"进度语义";进度的真相源回归到数据库自身。这比"队列 offset + 独立 checkpoint 表"健壮得多,因为原子性由单机事务保证,无分布式协调负担。
4.1 事务性批次写入 + 进度同库同事务
每批在一条事务内同时完成:① 数据幂等 upsert;② 推进 watermark(即断点)。提交成功 ⇔ 二者同时可见;事务中任何一步失败 ⇔ 整批回滚,库里无中间态。
4.2 预写式检查点(WAL 思想)
watermark 是预写日志的落地形态:它永远 ≤ 已提交数据的最大批次号。恢复时以 watermark 为准重放其后批次,天然无脏断点。
4.3 幂等键设计
以 url_hash 建唯一约束,ON CONFLICT DO NOTHING 让重放批次静默跳过,把 at-least-once 收敛为 effectively-once。幂等键必须覆盖"同一任务 + 同一 URL",跨任务则需加 job_id 维度。
4.4 网络层与存储层职责隔离
网络层失败(代理超时、连接 reset)由独立的代理/重试层吸收——例如通过隧道代理透明轮换出口 IP,其职责仅限于"让请求成功",绝不参与事务提交、绝不修改 watermark。跨层耦合是脏数据最常见的隐蔽来源。
五、生产级实战:Transactional Watermark Crawler
```import sqlite3
from contextlib import contextmanager
from dataclasses import dataclass
from typing import Iterable

@dataclass(frozen=True)
class PageRecord:
url: str
url_hash: str
payload: str

@contextmanager
def _txn(conn: sqlite3.Connection):
"""单一事务边界:异常回滚,正常提交。批量副作用的原子单元。"""
try:
conn.execute("BEGIN")
yield conn
conn.commit()
except Exception:
conn.rollback()
raise

class ExactlyOnceCrawler:
"""进度与数据同库同事务,从根上消除 offset/side-effect 鸿沟。"""

def __init__(self, conn: sqlite3.Connection, batch_size: int = 500) -> None:
    self.conn, self.batch_size = conn, batch_size

def _watermark(self) -> int:
    row = self.conn.execute(
        "SELECT watermark FROM crawl_watermark WHERE job_id=?", (JOB_ID,)
    ).fetchone()
    return row[0] if row else -1

def run(self, queue: Iterable[PageRecord]) -> None:
    buf: list[PageRecord] = []
    nxt = self._watermark() + 1
    for rec in queue:
        buf.append(rec)
        if len(buf) < self.batch_size:
            continue
        self._commit_batch(nxt, buf)
        buf, nxt = [], nxt + 1
    if buf:
        self._commit_batch(nxt, buf)

def _commit_batch(self, batch_id: int, batch: list[PageRecord]) -> None:
    with _txn(self.conn) as cur:
        cur.executemany(                      # ① 业务数据:幂等 upsert
            "INSERT INTO pages(url,url_hash,payload) VALUES (:url,:url_hash,:payload) "
            "ON CONFLICT(url_hash) DO NOTHING",
            [{"url": r.url, "url_hash": r.url_hash, "payload": r.payload} for r in batch],
        )
        cur.execute(                          # ② 进度:与数据同事务
            "INSERT INTO crawl_watermark(job_id,watermark) VALUES (?,?) "
            "ON CONFLICT(job_id) DO UPDATE SET watermark=excluded.watermark",
            (JOB_ID, batch_id),
        )
    # 提交成功 ⇔ 数据与进度同时落库,崩溃无中间态


多 Worker 扩展:乐观并发防止跳号
横向扩容时,多个 worker 可能竞争推进同一 watermark。用条件更新加乐观锁:
```cur.execute(
    "UPDATE crawl_watermark SET watermark=? WHERE job_id=? AND watermark<?",
    (batch_id, JOB_ID, batch_id),
)
if cur.rowcount == 0:
    raise ConcurrencyError("watermark 已被其他 worker 推进,本批丢弃重排")
WHERE watermark < ? 保证只有"推进"成功、"回退"被拒,从而避免批次跳号与漏采。

六、一致性验证
计数对账:SUM(每批写入行数) 应等于 pages 总行数,差为漏采。
唯一性约束:COUNT(*) - COUNT(DISTINCT url_hash) = 0,差为重复。
水位连续性:watermark 单调递增无跳号;跳号即某批被回滚却未被重调度。
校验和:对批次内 url_hash 集合取哈希入账,恢复后比对,快速定位错位批次。
建议将上述校验固化为恢复后的自动断言,不通过即告警而非继续。
七、小结
长周期爬虫的数据一致性,不靠"别崩"保证,而靠"崩了也能回到一致状态"保证。本文的模型是:
接受 at-least-once,以幂等键收敛为 effectively-once;
把进度变为数据,与业务写入置于同一 RDBMS 事务边界,使数据库成为唯一真相源;
隔离网络层与存储层,代理重试永不越界触碰事务。
铁律只有一条——watermark 永远 ≤ 已提交数据,任一受影响批次要么完整存在、要么根本不存在。 在此模型上,断点续爬、多 worker 扩容、代理容错都能各自演化而不破坏一致性。
示例为零依赖 SQLite 演示;生产环境将存储切换为 PostgreSQL/MySQL 连接池,并为 pages(url_hash) 与 crawl_watermark(job_id) 建立唯一约束以在库层强制幂等与单真相源。

相关文章
|
2月前
|
数据采集 Web App开发 JavaScript
全网电影信息爬取:从单机脚本到分布式采集系统的工程实践
全网电影信息爬取:从单机脚本到分布式采集系统的工程实践
|
3月前
|
数据采集 API 数据处理
Python 量化数据工程:基于 Pandas 的 A 股 K 线标准化清洗与向量化技术指标计算
Python 量化数据工程:基于 Pandas 的 A 股 K 线标准化清洗与向量化技术指标计算
|
3月前
|
JSON BI 调度
Python + 大模型行业资讯自动化摘要流水线完整工程实现方案
Python + 大模型行业资讯自动化摘要流水线完整工程实现方案
|
8月前
|
数据采集 JSON Java
Java 异步爬虫高效获取小红书短视频内容
Java 异步爬虫高效获取小红书短视频内容
|
9月前
|
数据采集 文字识别 JavaScript
基于文本检测的 Python 爬虫弹窗图片定位与拖动实现
基于文本检测的 Python 爬虫弹窗图片定位与拖动实现
|
3月前
|
数据采集 前端开发 JavaScript
Scrapling:极简高效的 Python 智能爬虫框架
Scrapling:极简高效的 Python 智能爬虫框架
|
5月前
|
数据采集 Web App开发 存储
Selenium+Python 爬虫:动态加载头条问答爬取
Selenium+Python 爬虫:动态加载头条问答爬取
Selenium+Python 爬虫:动态加载头条问答爬取
|
3月前
|
数据采集 存储 自然语言处理
信息筛选耗时?Python 爬虫搭配大模型,一键抓取资讯并智能总结
信息筛选耗时?Python 爬虫搭配大模型,一键抓取资讯并智能总结
|
4月前
|
数据采集 数据可视化 数据挖掘
均线选股策略研究:基于 Python 数据分析实现
均线选股策略研究:基于 Python 数据分析实现
|
4月前
|
数据采集 Web App开发 JavaScript
Python 爬虫动态 JS 渲染与无头浏览器实战选型指南
Python 爬虫动态 JS 渲染与无头浏览器实战选型指南

热门文章

最新文章