news 2026/9/20 13:29:11

ArchiveBox 耐用任务队列模型 ModelWithQueue 源码解析:status、retry_at 与原子租约机制

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
ArchiveBox 耐用任务队列模型 ModelWithQueue 源码解析:status、retry_at 与原子租约机制

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):任务的当前状态,默认QUEUEDnull=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 = 60
  • RETRY_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。

类级状态常量

类属性默认值说明
StatusChoicesDefaultStatusChoices状态枚举,子类可覆盖
INITIAL_STATEQUEUED初始状态
ACTIVE_STATESTARTED活跃状态
FINAL_STATES(SEALED,)终态集合
FINAL_OR_ACTIVE_STATES(*FINAL_STATES, ACTIVE_STATE)终态 + 活跃态
warn_on_save_outside_runnerTrue是否在 runner 进程外save()时告警

实例属性RETRY_ATSTATE只是retry_atstatus的属性别名,提供统一的读写接口。

生命周期操作: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 True
  • pause():将状态置为PAUSED并把retry_at推到RETRY_AT_MAX。三个保护条件——枚举未定义PAUSED、已是终态、已处于暂停——都会直接返回False。写库时使用safe_updateextra_filter={"status": previous_status},确保只有"仍处于原状态"的行才被更新(并发安全)。
  • resume(when=None):仅当当前is_paused时才生效,把状态改回QUEUEDretry_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 == 1
  • get_queue():取出所有"到期"(retry_at <= now)的行并按retry_at升序排列,构成待处理队列。这个查询天然排除了PAUSEDretry_at为远未来)和未来才到期的任务。
  • claim_for_worker():认领的核心,是一个**乐观锁 + CAS(比较并交换)**操作。更新条件同时包含pkretry_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 == 1
  • safe_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 = False

Binary 的模块 docstring 说明了设计意图:安装流程是同步的queued → installed转换,失败时保持queued并把retry_at顺延以便稍后重试;数据库行是唯一的生命周期状态,worker 中断后无需在内存中二次协调状态机。它还演示了update_and_requeue(retry_at=None, status=INSTALLED)的收尾用法(machine/models.py)。

运行器中的真实调用链

runner(services/runner.py)是把这套协议串起来的消费端,其处理每个到期任务的标准流程是:

  1. 选行:通过retry_at__lte=now类查询(等价于get_queue)取出到期任务。
  2. 认领:调用crawl.claim_processing_lock(lock_seconds=lock_seconds)(L1361、L1395)或Snapshot.claim_for_worker(snapshot, lock_seconds=lock_seconds)(L1448)原子获取租约;失败即返回False,说明该行已被其他 worker 认领。
  3. 执行与续租:处理期间用crawl.update_and_requeue(status=STARTED, retry_at=now + timedelta(seconds=ACTIVE_STATE_LEASE_SECONDS))续租(L1352),把租约从 60 秒起步按需延长。
  4. 收尾:终态任务(如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/resumeRETRY_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/20 13:26:04

基于Python的兵棋推演游戏源码解析与二次开发指南

简介&#xff1a;这是一份基于Python实现的兵棋推演游戏源码&#xff0c;面向对人工智能与战略模拟感兴趣的开发者&#xff0c;可用于学习智能体通信、指令处理与可视化推演流程。资源共35个文件&#xff0c;包括33个Python脚本、1个txt及1个markdown说明&#xff0c;压缩包仅1…

作者头像 李华
网站建设 2026/9/20 13:25:15

Phoenix 项目 Elixir 编码规范实战指南:从代码风格到高可靠测试

Phoenix 项目 Elixir 编码规范实战指南&#xff1a;从代码风格到高可靠测试 【免费下载链接】phoenix Peace of mind from prototype to production 项目地址: https://gitcode.com/gh_mirrors/ph/phoenix 本篇技术指南基于 Phoenix 框架仓库的 usage-rules/elixir.md 编…

作者头像 李华
网站建设 2026/9/20 13:24:45

OpenClaw 请求 401?TaoToken 这样核对 API 地址

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华