ArchiveBox 耐用任务队列模型 ModelWithQueue 源码解析:status、retry_at 与原子租约机制
【免费下载链接】ArchiveBox🗃 Open source self-hosted web archiving. Takes URLs/browser history/bookmarks/Pocket/Pinboard/etc., saves HTML, JS, PDFs, media, and more...项目地址: https://gitcode.com/gh_mirrors/ar/ArchiveBox
本文以 ArchiveBox 的archivebox.workers.models模块(API 文档见 archivebox.workers.models.md,实现见 workers/models.py)为核心,深入讲解 ArchiveBox 调度系统赖以运转的耐用队列(durable queue)协议:status/retry_at双字段状态模型、pause/resume生命周期控制、基于条件更新的原子租约(lease)与安全更新机制,以及它们在 Crawl、Snapshot、Binary 等具体模型和 runner 运行器中的实际调用方式。读完本文,你将掌握 ArchiveBox 分布式/多进程调度中"谁拥有任务、任务何时重试、如何防止重复执行"的底层实现原理,并了解如何在自己的 Django 模型上复用这套队列协议。
模块定位:调度系统的公共队列协议
ArchiveBox 使用 Django ORM 作为任务调度的中心化状态存储:爬虫任务(Crawl)、页面快照(Snapshot)、依赖二进制安装(Binary)都围绕"队列行"运转,由 runner.py 中的运行器循环消费。archivebox.workers.models就是这套体系的最小公共部分——它只定义了一个抽象基类ModelWithQueue,其模块 docstring 明确写道:
Durable queue fields and atomic lease operations shared by work rows. Concrete models own lifecycle transitions. This mixin only owns the common database queue protocol: status, retry_at, pause/resume, and claims.
也就是说:具体模型自己决定生命周期状态迁移规则,ModelWithQueue只负责提供通用的数据库队列协议——状态字段、重试时间、暂停/恢复、以及"认领"(claim)操作。这种"mixin 管协议、具体模型管迁移"的职责划分,是理解整个 workers 模块的关键。
队列状态协议:DefaultStatusChoices
模块首先定义了默认的状态枚举:
class DefaultStatusChoices(models.TextChoices): QUEUED = "queued", "Queued" STARTED = "started", "Started" PAUSED = "paused", "Paused" SEALED = "sealed", "Sealed"四个状态的含义如下:
| 状态 | 枚举值 | 含义 |
|---|---|---|
QUEUED | "queued" | 任务已入队,等待 worker 认领 |
STARTED | "started" | 任务正在处理中(活跃状态) |
PAUSED | "paused" | 任务被用户暂停,不再调度 |
SEALED | "sealed" | 任务已封存(终态,不再处理) |
DefaultStatusChoices继承自 Django 的models.TextChoices,因此同时具备枚举常量(DefaultStatusChoices.QUEUED)与字段 choices(.choices)两套能力。需要说明的是,这只是"默认"协议——使用方可以通过自定义StatusChoices完全重定义状态集(见下文 Binary 模型的定制),因此这四态是调度系统的常规约定而非铁律。
两个核心队列字段与默认值工厂
ModelWithQueue只有两个自有数据库字段,它们构成了队列的"双支柱":
default_status_field: models.CharField = models.CharField( choices=DefaultStatusChoices.choices, max_length=15, default=DefaultStatusChoices.QUEUED, null=False, blank=False, db_index=True, ) default_retry_at_field: models.DateTimeField = models.DateTimeField(default=timezone.now, null=True, blank=True, db_index=True)status(CharField,max_length=15):任务的当前状态,默认QUEUED,null=False, blank=False强制非空,且db_index=True保证按状态过滤的高效性。retry_at(DateTimeField,可空,默认当前时间):任务"何时到期可被认领"。它同时承担了到期时间与租约锁双重职责(下文详述),同样带数据库索引。
字段本身通过deconstruct()参数以模块级默认值工厂的形式定义(default_status_field/default_retry_at_field),方便各模型复用。ModelWithQueue内部据此声明字段:
status: models.CharField = models.CharField(**default_status_field.deconstruct()[3]) retry_at: models.DateTimeField = models.DateTimeField(**default_retry_at_field.deconstruct()[3])使用方可以在此基础上覆盖参数,例如 Binary 模型调用ModelWithQueue.StatusField(choices=StatusChoices.choices, default=StatusChoices.QUEUED, max_length=16)(见 machine/models.py)微调max_length。
两个关键常量
RETRY_AT_MAX = datetime(9999, 1, 1, tzinfo=UTC) ACTIVE_STATE_LEASE_SECONDS = 60RETRY_AT_MAX:一个"远未来"时间戳,语义上等价于"永不重试"。pause()会把retry_at置为该值,从而让被暂停的任务在get_queue()查询中永远不满足retry_at__lte=now条件,实现"不可见即不调度"。ACTIVE_STATE_LEASE_SECONDS = 60:活跃状态的默认租约时长。当 Crawl 保持STARTED状态继续推进子任务时,runner 会用它刷新父行的租约(见 runner.py),避免父行租约过期被其他 worker 抢走。
模块还导出了logger(模块日志器)、MODULE_PATH/REPO_ROOT/PACKAGE_ROOT(用于save()覆盖中定位调用方栈帧,见下文),这些是支撑调试告警的辅助常量。
ModelWithQueue:抽象基类的完整 API
ModelWithQueue继承django.db.models.Model,其Meta声明app_label = "workers"且abstract = True——它不产生数据表,只为子类提供字段与方法(从仓库看 workers/migrations 目录仅有__init__.py,正因抽象模型不生成迁移)。下面按职责分组介绍其完整 API。
类级状态常量
| 类属性 | 默认值 | 说明 |
|---|---|---|
StatusChoices | DefaultStatusChoices | 状态枚举,子类可覆盖 |
INITIAL_STATE | QUEUED | 初始状态 |
ACTIVE_STATE | STARTED | 活跃状态 |
FINAL_STATES | (SEALED,) | 终态集合 |
FINAL_OR_ACTIVE_STATES | (*FINAL_STATES, ACTIVE_STATE) | 终态 + 活跃态 |
warn_on_save_outside_runner | True | 是否在 runner 进程外save()时告警 |
实例属性RETRY_AT与STATE只是retry_at、status的属性别名,提供统一的读写接口。
生命周期操作:pause / resume / bump_retry_at
def pause(self, *, save: bool = True) -> bool: paused_state = getattr(self.StatusChoices, "PAUSED", None) if paused_state is None or self.status in self.FINAL_STATES or self.is_paused: return False previous_status = self.status self.status = paused_state self.retry_at = RETRY_AT_MAX if save: return self.safe_update( {"status": paused_state, "retry_at": RETRY_AT_MAX}, extra_filter={"status": previous_status}, ) return Truepause():将状态置为PAUSED并把retry_at推到RETRY_AT_MAX。三个保护条件——枚举未定义PAUSED、已是终态、已处于暂停——都会直接返回False。写库时使用safe_update且extra_filter={"status": previous_status},确保只有"仍处于原状态"的行才被更新(并发安全)。resume(when=None):仅当当前is_paused时才生效,把状态改回QUEUED,retry_at设为when(默认timezone.now(),即立即恢复调度),同样带extra_filter={"status": paused_state}的守卫。调用方可以指定when实现"定时恢复"。bump_retry_at(seconds=10):把retry_at从当前时刻顺延 N 秒,用于简单的指数退避或失败重试。is_paused:判断当前状态是否为枚举中定义的PAUSED。
队列查询与原子租约:get_queue / claim_for_worker / claim_processing_lock
这是模块最核心的并发机制,全部围绕retry_at实现"谁先到期谁被处理、谁先更新谁拥有":
@classmethod def get_queue(cls): return cls.objects.filter(retry_at__lte=timezone.now()).order_by("retry_at") @classmethod def claim_for_worker(cls, obj: "ModelWithQueue", lock_seconds: int = 60) -> bool: now = timezone.now() lock_until = now + timedelta(seconds=lock_seconds) updated = cls.objects.filter(pk=obj.pk, retry_at=obj.retry_at, retry_at__lte=now).update( retry_at=lock_until, modified_at=now, ) if updated == 1: obj.retry_at = lock_until cast(Any, obj).modified_at = now return updated == 1get_queue():取出所有"到期"(retry_at <= now)的行并按retry_at升序排列,构成待处理队列。这个查询天然排除了PAUSED(retry_at为远未来)和未来才到期的任务。claim_for_worker():认领的核心,是一个**乐观锁 + CAS(比较并交换)**操作。更新条件同时包含pk、retry_at=obj.retry_at(读到的旧值)和retry_at__lte=now(未过期)。只有恰好更新 1 行(updated == 1)才算认领成功,并把retry_at推进到"当前时间 + lock_seconds",相当于租约到期时间。两个并发 worker 同时认领同一行时,后执行的 UPDATE 因retry_at已被前一个改成租约时间而不满足retry_at=obj.retry_at条件,更新 0 行即失败——这就是"原子租约"防重入的原理。claim_processing_lock(lock_seconds=60):实例方法版认领。额外检查终态与retry_at is None,通过后委托给claim_for_worker。runner 处理每个任务前都调用它,拿到锁才继续执行。
安全更新:safe_update / update_and_requeue
def safe_update(self, update_fields, *, refresh=True, extra_filter=None) -> bool: values = dict(update_fields) values.setdefault("modified_at", timezone.now()) queryset = type(self).objects.filter(pk=self.pk) if extra_filter: queryset = queryset.filter(**extra_filter) updated = queryset.update(**values) if updated != 1 and extra_filter: current = type(self).objects.filter(pk=self.pk).values("status").first() logger.info( "SafeUpdateGuardMiss: %s row %s extra_filter=%s current_status=%s loaded_status=%s ...", ... ) if refresh: try: self.refresh_from_db() except type(self).DoesNotExist: pass return updated == 1safe_update():用带可选extra_filter的 UPDATE 实现"带守卫的原子更新",返回是否恰好更新 1 行。守卫失败时(如extra_filter指定的状态已被他人改变)会记录SafeUpdateGuardMiss日志,便于排查竞态。默认在更新后refresh_from_db()同步内存态。update_and_requeue(**kwargs):safe_update的特化版本,固定extra_filter={"retry_at": self.retry_at}——即"只有租约未变时才能更新",再次确认所有权后才允许改写状态。runner 中大量使用它推进 Crawl/Snapshot 的调度状态(见 runner.py)。
save() 覆盖:runner 进程外的保存告警
ModelWithQueue.save()覆写了 Django 的save(),实现了一个运行期安全网:当非新增行在非 ORCHESTRATOR(编排器)进程中被save()时,会通过栈帧回溯(利用MODULE_PATH/PACKAGE_ROOT/REPO_ROOT定位调用者,跳过本模块与 site-packages 帧)记录警告日志,指明调用方文件与行号。其意图是:队列行的状态变更应当由 runner 统一管理,直接在任意代码里save()可能绕过租约守卫造成状态错乱;该告警帮助开发者尽早发现这类"越权写库"。子类可通过warn_on_save_outside_runner = False关闭(Binary 模型即如此,因为它由BinaryService同步驱动安装流程,见 machine/models.py)。
辅助方法:status_counts / extend_choices / StatusField / RetryAtField
status_counts(queryset=None, statuses=None) -> dict[str, int]:类方法,按状态统计行数,返回如{"queued": 12, "started": 3, ...},用于调度概览与 admin 面板展示。extend_choices(base_choices):一个类装饰器工厂,把基础枚举与附加枚举合并生成新的TextChoices,供需要扩展状态集的模型使用。StatusField(**kwargs)/RetryAtField(**kwargs):类方法字段工厂,以默认字段定义为基础合并覆盖参数,子类声明字段时直接复用(Crawl、Snapshot、Binary 三处均如此)。
子类化实践:Crawl、Snapshot 与 Binary
从仓库看,ModelWithQueue被三个模型继承,展示了"默认协议复用 + 按需定制"的完整模式:
Crawl 与 Snapshot:直接复用默认四态
Crawl 与 Snapshot 都按默认协议定制:
# crawls/models.py class Crawl(ModelWithDeleteAfter, ModelWithOutputDir, ModelWithConfig, ModelWithHealthStats, ModelWithQueue): status = ModelWithQueue.StatusField( choices=ModelWithQueue.StatusChoices, default=ModelWithQueue.StatusChoices.QUEUED, ) retry_at = ModelWithQueue.RetryAtField(default=timezone.now) StatusChoices = ModelWithQueue.StatusChoices INITIAL_STATE = StatusChoices.QUEUED ACTIVE_STATE = StatusChoices.STARTED FINAL_STATES = (StatusChoices.SEALED,) FINAL_OR_ACTIVE_STATES = (*FINAL_STATES, ACTIVE_STATE)Snapshot 结构相同(core/models.py),并额外定义RUNNABLE_STATES = (QUEUED, STARTED)、OPEN_STATES等派生集合。两者共享默认四态:queued → started → sealed,可被paused中断。
Binary:完全重定义状态机
Binary 演示了最激进的定制——它只用两态:
class StatusChoices(models.TextChoices): QUEUED = "queued", "Queued" INSTALLED = "installed", "Installed" status = ModelWithQueue.StatusField(choices=StatusChoices.choices, default=StatusChoices.QUEUED, max_length=16) INITIAL_STATE = StatusChoices.QUEUED ACTIVE_STATE = StatusChoices.QUEUED # 活跃态即 QUEUED FINAL_STATES = (StatusChoices.INSTALLED,) warn_on_save_outside_runner = FalseBinary 的模块 docstring 说明了设计意图:安装流程是同步的queued → installed转换,失败时保持queued并把retry_at顺延以便稍后重试;数据库行是唯一的生命周期状态,worker 中断后无需在内存中二次协调状态机。它还演示了update_and_requeue(retry_at=None, status=INSTALLED)的收尾用法(machine/models.py)。
运行器中的真实调用链
runner(services/runner.py)是把这套协议串起来的消费端,其处理每个到期任务的标准流程是:
- 选行:通过
retry_at__lte=now类查询(等价于get_queue)取出到期任务。 - 认领:调用
crawl.claim_processing_lock(lock_seconds=lock_seconds)(L1361、L1395)或Snapshot.claim_for_worker(snapshot, lock_seconds=lock_seconds)(L1448)原子获取租约;失败即返回False,说明该行已被其他 worker 认领。 - 执行与续租:处理期间用
crawl.update_and_requeue(status=STARTED, retry_at=now + timedelta(seconds=ACTIVE_STATE_LEASE_SECONDS))续租(L1352),把租约从 60 秒起步按需延长。 - 收尾:终态任务(如
SEALED的 Crawl)在claim_for_worker成功后执行cleanup_runtime(),最后update_and_requeue(retry_at=None)清空调度时间,使其退出队列(L1405-L1411)。
值得注意的边界处理:当 Snapshot 的租约已过期但关联的abx-dl钩子进程仍在运行(process_set.filter(status="running"))时,runner 不会启动第二波处理,而是update_and_requeue(retry_at=now + lock_seconds)重新排队(L1471-L1479),保留快照级的所有权边界,防止重复执行。
测试验证与可观测性
仓库测试对这套协议有直接覆盖,可作为理解佐证:
- test_crawl_runner.py 断言
snapshot.claim_processing_lock(lock_seconds=60) is True,验证单 worker 场景下的认领成功路径; - 同文件 L414 使用
crawl.update_and_requeue(status=QUEUED, retry_at=timezone.now())构造调度场景,验证重排队逻辑。
此外,认领失败与守卫未命中都会留下可观测日志:SafeUpdateGuardMiss(safe_update 守卫失败)以及save() outside runner process告警,配合status_counts()的状态分布统计,可以在运维层面快速定位"任务堆积在哪一态、是否被越权改写"。
小结
archivebox.workers.models用约 230 行代码给出了一个干净、可复用的数据库级任务队列协议:以status+retry_at两个索引字段为存储,以"条件 UPDATE + 影响行数判定"实现原子认领与安全更新,以pause/resume与RETRY_AT_MAX实现暂停语义,以ACTIVE_STATE_LEASE_SECONDS维持活跃任务租约。Crawl、Snapshot、Binary 通过继承与覆盖各取所需,runner 则作为唯一合法的状态迁移入口消费队列——这套"mixin 管协议、模型管迁移、runner 管执行"的设计,是 ArchiveBox 多进程调度一致性的基石。
【免费下载链接】ArchiveBox🗃 Open source self-hosted web archiving. Takes URLs/browser history/bookmarks/Pocket/Pinboard/etc., saves HTML, JS, PDFs, media, and more...项目地址: https://gitcode.com/gh_mirrors/ar/ArchiveBox
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考