ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

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

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四个状态的含义如下状态枚举值含义QUEUEDqueued任务已入队等待 worker 认领STARTEDstarted任务正在处理中活跃状态PAUSEDpaused任务被用户暂停不再调度SEALEDsealed任务已封存终态不再处理DefaultStatusChoices继承自 Django 的models.TextChoices因此同时具备枚举常量DefaultStatusChoices.QUEUED与字段 choices.choices两套能力。需要说明的是这只是默认协议——使用方可以通过自定义StatusChoices完全重定义状态集见下文 Binary 模型的定制因此这四态是调度系统的常规约定而非铁律。两个核心队列字段与默认值工厂ModelWithQueue只有两个自有数据库字段它们构成了队列的双支柱default_status_field: models.CharField models.CharField( choicesDefaultStatusChoices.choices, max_length15, defaultDefaultStatusChoices.QUEUED, nullFalse, blankFalse, db_indexTrue, ) default_retry_at_field: models.DateTimeField models.DateTimeField(defaulttimezone.now, nullTrue, blankTrue, db_indexTrue)statusCharFieldmax_length15任务的当前状态默认QUEUEDnullFalse, blankFalse强制非空且db_indexTrue保证按状态过滤的高效性。retry_atDateTimeField可空默认当前时间任务何时到期可被认领。它同时承担了到期时间与租约锁双重职责下文详述同样带数据库索引。字段本身通过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(choicesStatusChoices.choices, defaultStatusChoices.QUEUED, max_length16)见 machine/models.py微调max_length。两个关键常量RETRY_AT_MAX datetime(9999, 1, 1, tzinfoUTC) ACTIVE_STATE_LEASE_SECONDS 60RETRY_AT_MAX一个远未来时间戳语义上等价于永不重试。pause()会把retry_at置为该值从而让被暂停的任务在get_queue()查询中永远不满足retry_at__ltenow条件实现不可见即不调度。ACTIVE_STATE_LEASE_SECONDS 60活跃状态的默认租约时长。当 Crawl 保持STARTED状态继续推进子任务时runner 会用它刷新父行的租约见 runner.py避免父行租约过期被其他 worker 抢走。模块还导出了logger模块日志器、MODULE_PATH/REPO_ROOT/PACKAGE_ROOT用于save()覆盖中定位调用方栈帧见下文这些是支撑调试告警的辅助常量。ModelWithQueue抽象基类的完整 APIModelWithQueue继承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_AT与STATE只是retry_at、status的属性别名提供统一的读写接口。生命周期操作pause / resume / bump_retry_atdef 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(whenNone)仅当当前is_paused时才生效把状态改回QUEUEDretry_at设为when默认timezone.now()即立即恢复调度同样带extra_filter{status: paused_state}的守卫。调用方可以指定when实现定时恢复。bump_retry_at(seconds10)把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__ltetimezone.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(secondslock_seconds) updated cls.objects.filter(pkobj.pk, retry_atobj.retry_at, retry_at__ltenow).update( retry_atlock_until, modified_atnow, ) if updated 1: obj.retry_at lock_until cast(Any, obj).modified_at now return updated 1get_queue()取出所有到期retry_at now的行并按retry_at升序排列构成待处理队列。这个查询天然排除了PAUSEDretry_at为远未来和未来才到期的任务。claim_for_worker()认领的核心是一个**乐观锁 CAS比较并交换**操作。更新条件同时包含pk、retry_atobj.retry_at读到的旧值和retry_at__ltenow未过期。只有恰好更新 1 行updated 1才算认领成功并把retry_at推进到当前时间 lock_seconds相当于租约到期时间。两个并发 worker 同时认领同一行时后执行的 UPDATE 因retry_at已被前一个改成租约时间而不满足retry_atobj.retry_at条件更新 0 行即失败——这就是原子租约防重入的原理。claim_processing_lock(lock_seconds60)实例方法版认领。额外检查终态与retry_at is None通过后委托给claim_for_worker。runner 处理每个任务前都调用它拿到锁才继续执行。安全更新safe_update / update_and_requeuedef safe_update(self, update_fields, *, refreshTrue, extra_filterNone) - bool: values dict(update_fields) values.setdefault(modified_at, timezone.now()) queryset type(self).objects.filter(pkself.pk) if extra_filter: queryset queryset.filter(**extra_filter) updated queryset.update(**values) if updated ! 1 and extra_filter: current type(self).objects.filter(pkself.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 / RetryAtFieldstatus_counts(querysetNone, statusesNone) - 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( choicesModelWithQueue.StatusChoices, defaultModelWithQueue.StatusChoices.QUEUED, ) retry_at ModelWithQueue.RetryAtField(defaulttimezone.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(choicesStatusChoices.choices, defaultStatusChoices.QUEUED, max_length16) 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_atNone, statusINSTALLED)的收尾用法machine/models.py。运行器中的真实调用链runnerservices/runner.py是把这套协议串起来的消费端其处理每个到期任务的标准流程是选行通过retry_at__ltenow类查询等价于get_queue取出到期任务。认领调用crawl.claim_processing_lock(lock_secondslock_seconds)L1361、L1395或Snapshot.claim_for_worker(snapshot, lock_secondslock_seconds)L1448原子获取租约失败即返回False说明该行已被其他 worker 认领。执行与续租处理期间用crawl.update_and_requeue(statusSTARTED, retry_atnow timedelta(secondsACTIVE_STATE_LEASE_SECONDS))续租L1352把租约从 60 秒起步按需延长。收尾终态任务如SEALED的 Crawl在claim_for_worker成功后执行cleanup_runtime()最后update_and_requeue(retry_atNone)清空调度时间使其退出队列L1405-L1411。值得注意的边界处理当 Snapshot 的租约已过期但关联的abx-dl钩子进程仍在运行process_set.filter(statusrunning)时runner 不会启动第二波处理而是update_and_requeue(retry_atnow lock_seconds)重新排队L1471-L1479保留快照级的所有权边界防止重复执行。测试验证与可观测性仓库测试对这套协议有直接覆盖可作为理解佐证test_crawl_runner.py 断言snapshot.claim_processing_lock(lock_seconds60) is True验证单 worker 场景下的认领成功路径同文件 L414 使用crawl.update_and_requeue(statusQUEUED, retry_attimezone.now())构造调度场景验证重排队逻辑。此外认领失败与守卫未命中都会留下可观测日志SafeUpdateGuardMisssafe_update 守卫失败以及save() outside runner process告警配合status_counts()的状态分布统计可以在运维层面快速定位任务堆积在哪一态、是否被越权改写。小结archivebox.workers.models用约 230 行代码给出了一个干净、可复用的数据库级任务队列协议以statusretry_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),仅供参考
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表