数据抓取

设计幂等爬虫:中断后重启而不重复

爬虫会崩溃。我们在抓取过程中将一个爬虫强制终止了十次:简单版本存储的行数多达4.3倍,幂等版本则恰好是500行。如何构建后者。

Chris Collins

Chris Collins

2026年10月2日 · 4 分钟阅读

每个爬虫最终都会在中途停下来。容器被重新调度,部署重启了工作进程,机器内存耗尽,或者有人按下了 Ctrl+C。接下来发生的事情决定了你的数据是否还值得信赖。一个只是简单地从头开始重跑的爬虫会重新抓取所有内容,而且除非专门设计过,还会把所有内容再写入一遍。

幂等爬虫指的是:同一份工作运行两次与运行一次效果相同。它可以在任何时刻被杀死、重启,最终每条记录仍然只有一份正确的副本。本指南展示了实现这一点的四个设计决策,附有经过测试的代码,以及在爬取过程中杀死它十次的结果。

关键要点

  • 假设进程会在任意两行代码之间死掉。设计时要确保重启最多只会重复一个正在处理中的页面。
  • 为每条记录生成一个源自其本身内容的 ID,比如来源和 SKU,而不是源自抓取时间。这样重复写入就是更新,而不是重复记录。
  • 将爬取的前沿队列保存在持久化存储中,并在存储结果的同一个事务中把 URL 标记为完成。
  • 使用租约而不是锁,这样被崩溃的工作进程占用的任务会自动重新变为可用状态。
  • 在我们对 500 个产品的测试中,每次运行杀死十次:朴素爬虫最终得到 1,286 到 2,165 行,请求数最多达到 4.3 倍;幂等爬虫最终得到恰好 500 行正确数据,重复请求最多 3 次。

为什么重启会产生重复数据

爬虫常见的最初版本是这样的:从第一页开始,跟随分页,为每个产品插入一行,边抓取边提交。它运行完美,直到被中断。重启时它完全不记得自己抓到了哪里,于是从头再来,每个已经存储过的产品都会以一个新的自增 ID 被再次插入。杀死它几次之后,表里大部分产品都有好几份副本,而且没有任何标记能说明哪一份是最新的。

事后去重是可行的,实体解析这篇文章讨论过这个方法,但那只是在处理症状。重复数据在影响准确性之前先造成了成本损失:每个重复抓取的页面都要支付两次带宽费用,这直接影响到每条干净记录的成本。

让爬虫幂等的四个决策

1. 从数据本身派生记录 ID

一条记录的 ID 应该来自这条记录本身是什么:来源加上它的自然键,比如 SKU、列表 ID 或规范化的 URL。把它们哈希在一起,同一个产品在每次运行、每台机器上都会得到相同的 ID。再次写入就变成了 upsert(更新插入):如果记录已存在且未变化,什么都不做;如果变化了,就更新它并记录变化时间。

2. 保持前沿队列持久化

待访问的 URL 列表,以及哪些已经完成,应该保存在数据库里,而不是内存中。添加一个已知的 URL 必须是无操作的,这样在已处理过的列表页上重新发现链接就不会造成影响。

3. 把结果和进度一起提交

危险的时刻发生在存储结果和记录该 URL 已完成之间。如果进程在完成一个之后、另一个之前死掉,重启要么会丢失结果,要么会重复它。在一个数据库事务中同时完成这两件事,就消除了这个间隙:要么记录被存储且 URL 被标记完成,要么两者都没发生,URL 会被简单地重新抓取。

4. 租用任务而不是锁定任务

工作进程在有限的时间内认领一个 URL。如果它完成了,URL 就被标记为完成。如果它崩溃了,租约到期,另一个工作进程,或者重启后的同一个进程,就会接手这个 URL。不需要任何人注意到崩溃,工作就能自动恢复。

代码

下面的模块使用 SQLite 和 requests 实现了上述全部四点。SQLite 让这个示例保持自包含;同样的设计可以迁移到 PostgreSQL 或任何支持事务和 upsert 的数据库。

import hashlib
import json
import re
import sqlite3
import time

import requests

SCHEMA = """
CREATE TABLE IF NOT EXISTS frontier (
    url TEXT PRIMARY KEY,
    status TEXT NOT NULL DEFAULT 'pending',      -- pending, leased, done
    leased_until REAL NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS records (
    record_id TEXT PRIMARY KEY,                  -- derived from the source and its natural key
    url TEXT NOT NULL,
    content_hash TEXT NOT NULL,
    data TEXT NOT NULL,
    first_seen REAL NOT NULL,
    last_changed REAL NOT NULL
);
"""


def open_db(path):
    db = sqlite3.connect(path, isolation_level=None)  # explicit transactions below
    db.execute("PRAGMA journal_mode=WAL")
    db.executescript(SCHEMA)
    return db


def record_id(source, natural_key):
    """The same item always gets the same ID, however many times it is scraped."""
    return hashlib.sha256(f"{source}:{natural_key}".encode()).hexdigest()[:16]


def enqueue(db, urls):
    """Adding a URL that is already known is a no-op, so re-discovering links is harmless."""
    db.executemany("INSERT OR IGNORE INTO frontier (url) VALUES (?)", [(u,) for u in urls])


def lease(db, lease_seconds=60):
    """Claim one URL. A lease left behind by a crashed worker expires and the URL is retried."""
    now = time.time()
    row = db.execute(
        "UPDATE frontier SET status = 'leased', leased_until = ? WHERE url = ("
        "  SELECT url FROM frontier WHERE status = 'pending' OR (status = 'leased' AND leased_until < ?) LIMIT 1"
        ") RETURNING url", (now + lease_seconds, now)).fetchone()
    return row[0] if row else None


def complete(db, url, new_urls=(), record=None):
    """Write the result and mark the URL done in one transaction: either both happen or neither does."""
    db.execute("BEGIN IMMEDIATE")
    try:
        enqueue(db, new_urls)
        if record:
            now = time.time()
            payload = json.dumps(record["data"], sort_keys=True)
            digest = hashlib.sha256(payload.encode()).hexdigest()
            db.execute(
                "INSERT INTO records (record_id, url, content_hash, data, first_seen, last_changed) VALUES (?, ?, ?, ?, ?, ?) "
                "ON CONFLICT(record_id) DO UPDATE SET data = excluded.data, content_hash = excluded.content_hash, "
                "last_changed = excluded.last_changed WHERE records.content_hash != excluded.content_hash",
                (record["id"], url, digest, payload, now, now))
        db.execute("UPDATE frontier SET status = 'done' WHERE url = ?", (url,))
        db.execute("COMMIT")
    except Exception:
        db.execute("ROLLBACK")
        raise


def crawl(db, base, session=None, lease_seconds=60):
    session = session or requests.Session()
    enqueue(db, [f"{base}/list/0"])
    while True:
        url = lease(db, lease_seconds)
        if url is None:
            if db.execute("SELECT 1 FROM frontier WHERE status != 'done' LIMIT 1").fetchone():
                time.sleep(1)  # a lease is still held, perhaps by a worker that crashed; wait for it to expire
                continue
            return
        html = session.get(url, timeout=30).text
        if "/list/" in url:
            links = [base + href for href in re.findall(r'href="(/(?:product|list)/\d+)"', html)]
            complete(db, url, new_urls=links)
        else:
            sku = re.search(r'class="sku">([^<]+)<', html).group(1)
            price = re.search(r'class="price">([^<]+)<', html).group(1)
            complete(db, url, record={"id": record_id("example-shop", sku), "data": {"sku": sku, "price": price}})

crawl 函数是针对我们的测试站点设计的,该站点有 25 个列表页面,链接到 500 个产品页面。它之上的所有内容都是可复用的。有两个细节容易被忽略:

  • 只有当内容发生变化时,upsert 才会写入。 冲突更新上的 WHERE 子句会让未变化的记录保持不动,这样 last_changed 才真正名副其实,可以用来驱动变化检测。
  • 循环不会在队列为空时停止。 我们的第一个版本在无法租用任何 URL 时就结束,这会导致每当崩溃留下一个未到期的租约时,就有一个页面未完成。现在的循环会一直等待,直到每个 URL 都完成为止。

杀死它十次

我们运行了一个本地测试站点,有 25 个列表页面和 500 个产品页面,所以一次完整爬取恰好需要 525 次请求。每次运行都启动爬虫,在 0.2 到 1 秒之间的随机时刻用 SIGKILL 杀死它,这样重复十次后再让它完成。我们用相同的杀死时间点,分别对上面的幂等爬虫和前面描述的朴素爬虫各运行五次,幂等版本使用 3 秒的租约。

运行次数朴素爬虫:存储行数朴素爬虫:请求数幂等爬虫:存储行数幂等爬虫:请求数
12,1652,283500525
21,8751,977500526
31,8621,967500526
41,4551,538500526
51,2861,361500528

两个版本最终都存储了全部 500 个产品,因为每次都完成了最后一次运行。区别在于其他一切。朴素爬虫每个产品存储了 2.6 到 4.3 行,发出的请求数是必要数量的 2.6 到 4.3 倍。幂等爬虫每个产品恰好存储一行,每个值都与源页面匹配,在十次崩溃中重复请求最多三次:对应每次恰好在页面处理中落下的杀死时刻。它在每次运行中也更快完成,耗时 6.8 到 8.3 秒,而朴素爬虫是 9.2 到 10.8 秒,尽管它有时需要等待一个崩溃工作进程的租约到期。

超越数据库

记录并不是爬虫唯一会重复执行的事情。同样的原则适用于每一个副作用:

  • 文件。 用内容的哈希值或记录 ID 来命名下载的图片和文档,这样重复下载会覆盖而不是新增一份副本。
  • 消息和 webhook。 发送一个源自记录及其版本的幂等键,这样消费者就可以忽略它已经处理过的消息。
  • 计数器和聚合值。 从存储的记录重新计算它们,而不是每次抓取时递增,否则重启会使它们虚增。
  • 原始响应。 如果你归档原始响应,重复抓取会增加一份额外的捕获记录,这是无害的,甚至有用,只要提取出的记录本身保持幂等。

实践准则

  • 仔细挑选自然键。 它在源站点上必须是稳定的:一个 SKU 或列表 ID,而不是页面上的位置或带有追踪参数的 URL。
  • 先让重试变得安全,再考虑频率。 一旦写入变得幂等,重试与退避就可以放心地频繁使用,而不会污染数据。
  • 根据最慢的页面设置租约时长。 比一次慢速抓取更短的租约会让两个工作进程处理同一个 URL;这里是无害的,但会浪费带宽。
  • 关注重复率。 每条存储记录对应的请求数是一个低成本的健康指标;它上升意味着出现了崩溃或租约问题,应该和其他指标一起纳入监控抓取管道。

总结

一个无法安全重启的爬虫最终会悄无声息地破坏自己的数据。解决办法不是更小心地操作,而是一种让重启变得无关紧要的设计:源自数据本身的 ID、持久化的前沿队列、结果与进度一起提交,以及会到期的租约。

在我们的测试中,这四个决策把十次崩溃造成的最多 4.3 倍的行数和请求数,变成了恰好正确的 500 条记录和三次重复请求。

来源与参考资料

  • SQLite 文档,UPSERT 和 RETURNING。
  • 该杀死测试由 Shifter 于 2026 年 10 月 2 日针对本地测试站点运行,使用上述代码。

准备好开始了吗?

试用 Shifter 住宅代理,205M+ 个 IP,195+ 个国家,低至 $0.10/GB。

立即开始