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