ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

CyberEdge 任务引擎原理解读:PostgreSQL FOR UPDATE SKIP LOCKED 高可靠任务调度完整指南

CyberEdge 任务引擎原理解读:PostgreSQL FOR UPDATE SKIP LOCKED 高可靠任务调度完整指南 CyberEdge 任务引擎原理解读PostgreSQL FOR UPDATE SKIP LOCKED 高可靠任务调度完整指南【免费下载链接】CyberEdge互联网资产综合扫描/攻击面测绘项目地址: https://gitcode.com/gh_mirrors/cy/CyberEdgeCyberEdge是一款 AI 原生的互联网资产综合扫描 / 攻击面测绘平台它的任务引擎没有引入任何独立消息队列而是基于PostgreSQL的FOR UPDATE SKIP LOCKED实现了高可靠的任务调度不丢任务、不重复执行、崩溃可恢复。本文带你完整看懂这套任务引擎的设计原理与关键实现。一、为什么用数据库当消息队列做任务调度时最常见的方案是引入 Redis、Kafka、RabbitMQ 等中间件。但 CyberEdge 在架构文档里给了一个明确判断任务用FOR UPDATE SKIP LOCKED领取状态变更、任务事件、发现记录、outbox 事件全部事务性提交。只有当实测吞吐量证明 PostgreSQL 不够用时才引入专用 broker。—— 出自 docs/architecture.md这个决策对自托管产品非常友好✅少一个组件部署时只需要一个 PostgreSQL不用维护额外队列✅强一致任务状态、事件记录、通知事件在同一个事务里提交不会出现任务完成了但事件没记上的裂脑情况✅天然持久化任务行落在磁盘里进程崩溃重启后队列原封不动二、任务引擎的整体结构 ️CyberEdge 由三个常驻协程组成入口在 src/main.rs角色轮询间隔职责调度器1 秒扫描到期计划批量生成任务入队工作器500 毫秒领取并执行资产扫描任务通知分发器1 秒从 outbox 领取待投递事件推送 webhook流程可以概括为计划(Schedule) 到期 → 生成任务(Task) 入队 → 工作器领取(claim) → 执行扫描 → 事务性落库 → outbox 事件投递其中计划由 migrations/0003_schedules.sql 定义的schedules表承载配合一个部分索引按到期时间排序任务则存在tasks表中核心是state状态字段。三、领取任务的关键FOR UPDATE SKIP LOCKED 原理 工作器每次轮询都会调用claim_task核心 SQL 只有一行见 src/repository/postgres.rs 第 761-807 行SELECT id FROM tasks WHERE state Queued ORDER BY created_at_seconds, created_at_nanos FOR UPDATE SKIP LOCKED LIMIT 1这行 SQL 是整个任务引擎的灵魂拆开看FOR UPDATE对选中的行加排他锁。同一事务内立刻把状态从Queued改为Running其他工作器看不到、动不了这行。SKIP LOCKED如果某个候选行正被其他事务锁住不等待直接跳过去取下一行。这避免了多个工作器同时抢同一行时互相阻塞的死锁式等待。ORDER BY created_atFIFO 公平性——最早入队的任务最先被领取。效果就是N 个工作器并发轮询时各自原子地拿到互不相同的任务队列为空时返回空结果安静退出没有任何锁竞争开销。四、一个事务锁住四件事 claim_task在同一个数据库事务里一次性完成① 把任务状态Queued → Running加锁领取② 插入task.running任务事件审计轨迹migrations 中的事件表③ 插入 outbox 事件用于后续 webhook 通知④ 读取任务与范围Scope并随结果返回任何一步失败整个事务回滚任务回到队列里等待下一次领取——要么全部生效要么像没发生过一样。这就是高可靠的第一道保障。任务执行完毕后complete_task再次校验任务必须处于Running状态ensure_running才允许写入发现记录并完成防止已失败/已取消的任务被错误收尾。五、崩溃与重复执行怎么防️故障场景引擎行为工作器进程崩溃任务停留在Running重启后可被审计/标记失败不会悄悄丢失两个工作器同时抢同一任务行锁 SKIP LOCKED保证只有一个事务能更新成功另一个直接跳过任务已Failed又被收尾ensure_running校验拦截写入被拒计划调度重复触发到期计划同样用FOR UPDATE SKIP LOCKEDsrc/repository/postgres.rs 第 366-441 行批量锁定后更新next_run_at并发下每个计划只会被推进一次状态机的流转是严格单向的Queued → Running → Completed / Failed每一次流转都伴随一条带序号sequence的任务事件写入事后可完整回放。六、通知也不丢outbox 模式 任务完成产生的通知不直接发 webhook而是先写入outbox_events表由独立分发器投递。领取通知用的还是同款技能... WHERE published_at IS NULL AND next_attempt_at now() AND (lease_until IS NULL OR lease_until now()) ORDER BY next_attempt_at, created_at FOR UPDATE SKIP LOCKED LIMIT 1额外的可靠性设计src/repository/postgres.rs 第 42-101 行30 秒租约lease_until防止分发器拿到事件后崩溃事件被永久占用指数退避重试第 n 次失败后等待2^n秒上限 900 秒死信兜底重试 8 次仍失败则标记 dead letter不再无限重试七、什么时候才需要引入独立消息队列CyberEdge 的态度写得很直白实测吞吐量证明不够用之前不加 broker。这套 PostgreSQL 任务调度足以支撑单组织自托管场景下的资产扫描任务当出现超大规模任务洪峰、需要跨机房扇出时再按 docs/ai-native-architecture.md 中的边界扩展。八、核心代码索引 模块路径任务领取 / 完成 / outbox核心src/repository/postgres.rs仓库抽象含内存实现便于测试src/repository/mod.rs、src/repository/memory.rs工作器主循环与策略分发src/worker.rs三协程启动入口src/main.rs计划表与到期索引migrations/0003_schedules.sql架构决策说明docs/architecture.md总结CyberEdge 的任务引擎证明了一个设计得当的 PostgreSQL完全可以替代消息队列承载任务调度。三根支柱——FOR UPDATE SKIP LOCKED实现无阻塞并发领取状态 事件 通知同事务提交保证原子性严格状态机 租约 指数退避兜住崩溃与失败如果你正在设计自己的资产扫描器或任务系统这套数据库即队列的方案值得直接借鉴。【免费下载链接】CyberEdge互联网资产综合扫描/攻击面测绘项目地址: https://gitcode.com/gh_mirrors/cy/CyberEdge创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表