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

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

一、重新定义问题:爬虫是一类有状态的流处理任务
长周期爬虫(任务跨度数天、采集千万级页面)在工程上不应被当成"脚本",而应被建模为一条有状态的流处理管道:待爬 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) 建立唯一约束以在库层强制幂等与单真相源。

相关文章
|
9天前
|
存储 弹性计算 缓存
阿里云服务器租赁费用:新版租赁收费标准及活动报价参考
本文更新了2026年阿里云全系列云服务器租赁活动报价,所有特惠资源均可前往阿里云活动中心选购,整体覆盖从个人入门到企业级高性能场景的全梯度需求。其中轻量应用服务器主打极致性价比,2核2G峰值200M带宽配置每日10点、15点限时抢购价仅38元/年,2核4G配置379元/年起;高性价比的经济型e实例、通用算力型u2i实例覆盖2核4G至4核32G全档位,适配开发测试与中小型企业业务;搭载英特尔至强6处理器的第九代c9i企业级实例算力较上代提升20%,支撑高并发生产环境,不同实例规格价差清晰,用户可根据自身业务负载与预算灵活选型。
1893 119
阿里云服务器租赁费用:新版租赁收费标准及活动报价参考
|
10天前
|
人工智能 程序员 API
Codex 接入 DeepSeek-V4-Flash:还能补上识图,提供两套方案
Codex 接入 DeepSeek-V4-Flash 怎么配?本文覆盖 CLI 与桌面端,再用 qwen3-vl-flash 补识图,两套方案可直接照做
1451 13
|
16天前
|
云安全 人工智能 运维
阿里云联动百位企业安全专家,共识Agent防御最佳实践
当Agent成为新员工,你的安全边界在哪里?
1966 10
阿里云联动百位企业安全专家,共识Agent防御最佳实践
|
7天前
|
编解码 弹性计算 云计算
MiniMax-H3 视频生成模型 — 一键部署与使用指南
MiniMax-H3是MiniMax开源的33B全模态视频生成模型,支持文生视频、图生视频、参考生视频三种模式,原生输出2K/15秒带立体声音频视频,已原生适配ComfyUI,并可通过阿里云计算巢一键部署。(239字)
|
10天前
|
人工智能 JSON Shell
2026AI漫剧本地全开源方案(附各个软件模型链接),8G显卡也能流畅运行
这是一套完全本地化部署的AI漫剧生成技术链路:涵盖LLM剧本分镜生成、FLUX文生图(IP-Adapter人脸锁定)、StoryDiffusion时序连贯控制、LTX-2.3唇形同步视频生成,及ComfyUI全流程调度。零云端费用,仅耗硬件算力,单集2–4小时可产出竖屏短视频,适配抖音/B站分发。
|
8天前
|
人工智能 API 开发工具
2026 零基础本地 AI 漫剧完整实操教程(8G 笔记本显卡可用|附可直接复制命令与代码)
本方案提供完全离线、本地运行的漫剧全自动制作流程:RTX3060/4050 8G显卡即可驱动,涵盖Qwen写分镜→ComfyUI统一角色绘图→LTX2.3图生微动画→Qwen3-TTS本地配音→FFmpeg自动合成,全程无水印、免API、不限次。专为低显存优化,解决变脸、闪烁、爆内存三大痛点。(239字)
|
22天前
|
人工智能 前端开发 Linux
Codex 桌面版安装 + CC Switch 接入第三方 API 完整教程(2026 最新)
2026最新教程:手把手教你安装Codex桌面版,通过CC Switch v3.17.0一键接入Fenno等国产API(兼容OpenAI Responses格式),跳过账号登录,完整启用代码审查、多步任务与上下文感知功能。零基础友好,全程图文实操。(239字)
3412 5
|
10天前
|
编解码 人工智能 安全
2核4G/4核8G/8核16G阿里云服务器如何选择实例?经济型e、通用算力型u2i与计算型c9i选哪个?
本文介绍了阿里云2核4G、4核8G、8核16G三档主流配置下经济型e、通用算力型u2i和计算型c9i三种实例的最新活动价格与适用场景。同配置下三者价差显著,以2核4G为例,经济型e低至599.93元/年,计算型c9i则高达1742.08元/年。文章详细解析了各实例的性能定位:经济型e适合轻负载入门场景,u2i兼顾稳定算力与性价比,c9i凭借第9代至强处理器与芯片级安全能力支撑高性能业务。同时提示用户可叠加满减优惠券享受折上折,建议根据业务负载与预算综合决策。
555 113