- Agent 框架
- 后端
- 低代码
- RAG
【免费下载链接】yao
✨ All your agents and workspaces in one place, on every device you own. Track tasks on a board, accessible from desktop, mobile, browser, or API. Self-hosted.
导读
Yao(当前仓库 job/ 目录)内置了一套完整的任务调度与执行框架 Job Framework,它在一个统一的模型体系下解决了任务创建、调度、执行、进度上报、日志记录与状态持久化的全链路问题。本文以 job/README.md 为主体,结合 job/types.go、job/worker.go、job/job.go、job/goroutine.go、job/process.go、job/data.go 等源码与 yao/models/job/ 下的数据模型定义,带你掌握一次性任务(Once)、Cron 定时任务与守护任务(Daemon)的创建方式,理解 Goroutine 与独立进程两种执行模式各自的适用场景,并学会使用进度追踪、结构化日志与完整的 CRUD API 构建可观测、可恢复的自托管任务系统。
一、框架概览:一个任务从创建到落库的完整闭环
Job Framework 是一套 "comprehensive task scheduling and execution framework",其设计目标可以概括为四个字:一切持久化。从任务的元数据、每次执行的实例记录、执行过程中的进度更新,到每一条日志,全部写入数据库(默认基于__yao.job、__yao.job.category、__yao.job.execution、__yao.job.log四个内置模型),因此系统重启后可以借助 job/job.go 中的RestoreJobsFromDatabase()从数据库恢复活跃任务,实现"定义即持久、恢复即继续"。
从源码结构看,框架按职责划分为五个层次(对应 job/README.md 的 Architecture 章节):
| 层次 | 文件 | 职责 |
|---|---|---|
| 数据层 | job/data.go | Job / Category / Execution / Log 四类实体的全部 CRUD 操作 |
| 执行层 | job/execution.go、job/job.go | 执行实例的创建、任务生命周期(Push/Stop/Destroy)管理 |
| Worker 层 | job/worker.go | Worker 池管理、任务分发与并发控制 |
| 进度层 | job/progress.go、job/interfaces.go | 进度追踪与数据库持久化 |
| 类型层 | job/types.go | 全部数据结构、枚举常量与注册表 |
二、数据模型:四个内置模型撑起任务全生命周期
数据模型定义位于 yao/models/job/ 目录,共四个.mod.yao文件:job.mod.yao、category.mod.yao、execution.mod.yao、log.mod.yao。它们与 job/types.go 中的Job、Category、Execution、Log结构体一一对应。
2.1 Job:任务主体
以 yao/models/job/job.mod.yao 为准,Job 表承载任务的完整元数据,关键字段如下:
- 身份与展示:
job_id(64 位唯一字符串标识)、name、icon、description、sort、created_by; - 调度配置:
schedule_type(枚举:once/cron/daemon,默认once)、schedule_expression(Cron 表达式,如"0 2 * * *")、next_run_at/last_run_at时间戳; - 执行配置:
mode(枚举:GOROUTINE/PROCESS,默认GOROUTINE)、max_worker_nums(默认 1)、max_retry_count(默认 0)、default_timeout(默认超时秒数,可被单次执行覆盖)、priority(默认 0,数值越大优先级越高); - 状态控制:
status(枚举 9 种:draft/ready/queued/running/paused/completed/failed/cancelled/disabled,默认draft)、enabled、system、readonly、current_execution_id; - 扩展字段:
config(JSON 配置)、Yao 平台字段(__yao_created_by、__yao_updated_by、__yao_team_id、__yao_tenant_id)。
模型还预置了四个复合索引以优化常见查询路径:idx_job_category_status(分类+状态)、idx_job_mode_status(模式+状态)、idx_job_schedule_type_next_run(调度类型+下次运行时间)、idx_job_created_by_status(创建人+状态),并开启了permission与timestamps选项。
2.2 Category:自动管理的分类体系
Category(job/types.go)用于组织任务。框架实现了"按名自动创建"的机制:在 job/data.go 的GetOrCreateCategory(name, description)中,先按name查询,不存在则自动创建;特殊地,名称为"Default"的分类会固定使用category_id = "default"并标记为system分类(见 job/data.go 的ensureCategoryExists)。这也解释了 job/job.go 中makeJob的行为:当CategoryID与CategoryName均为空时,任务默认归入"Default"分类。
2.3 Execution:每次执行的实例快照
Execution(job/types.go)是任务每次运行的实例记录,字段极为丰富,基本做到"一次执行一份审计档案":
- 触发信息:
trigger_category、trigger_source、trigger_context、scheduled_at; - 运行归属:
worker_id、process_id、parent_execution_id(支持父子执行链)、retry_attempt; - 时间与结果:
started_at、ended_at、timeout_seconds、duration(毫秒)、result、error_info、stack_trace、metrics、context; - 执行配置快照:
execution_config(运行时使用)与config_snapshot(JSON 落库快照)。job/data.go 在查询时会从ConfigSnapshot反序列化还原ExecutionConfig,确保即使进程重启,执行配置也能完整恢复。
2.4 Log:结构化执行日志
Log(job/types.go)记录任务执行过程中的每一条事件,字段包括level、message、context、source、execution_id、step、progress、duration、error_code、stack_trace、worker_id、process_id、timestamp、sequence(序列号)。job/data.go 的ListLogs默认按timestamp倒序分页查询,SaveLog则始终追加新记录。
三、两种执行模式:Goroutine 与独立进程
执行模式由ModeType常量定义(job/types.go):
const ( GOROUTINE ModeType = "GOROUTINE" // Execute using Go goroutine (lightweight, fast) PROCESS ModeType = "PROCESS" // Independent process isolated )两种模式在 job/worker.go 中通过switch w.Mode分流,最终分别调用executeInGoroutine与executeInProcess。Worker 池中所有 Worker 默认以GOROUTINE模式启动(见 job/worker.go)。
3.1 Goroutine 模式:轻量、快速、共享内存
在executeInGoroutine(job/worker.go)中,框架根据执行配置的类型(ExecutionType,见 job/types.go)分派到三种执行器:
| ExecutionType | 说明 | 实现 |
|---|---|---|
process | Yao 处理器(默认类型) | job/goroutine.goExecuteYaoProcess |
command | 系统命令 | job/goroutine.goExecuteSystemCommand |
func | Go 函数(内存注册表,内部使用) | job/goroutine.goExecuteFunc |
Goroutine 模式执行 Yao 处理器时,通过process.NewWithContext(ctx, ...)创建带上下文的进程对象,并利用proc.WithCallback(...)注册回调来接收 Yao 侧上报的进度数据(详见第五章)。共享数据(ExecutionOptions.SharedData)会通过proc.WithGlobal/proc.WithSID/proc.WithAuthorized注入进程上下文,从而把会话 ID、鉴权信息等一并传递到业务处理器中。
3.2 Process 模式:独立进程、彻底隔离
executeInProcess(job/worker.go)会生成proc_xxxxxxxx形式的process_id,并调用 job/process.go 中对应的执行器:
- 执行 Yao 处理器:通过
exec.CommandContext(ctx, "yao", args...)调用yao run <processName> <args>命令,工作目录设置为config.Conf.Root(Yao 应用根目录),并通过环境变量YAO_JOB_ID、YAO_EXECUTION_ID以及YAO_JOB_SHARED_<key>(共享数据序列化为 JSON 后注入)向子进程传递上下文(job/process.go); - 执行系统命令:同样以独立子进程方式运行
exec.CommandContext,并注入相同的环境变量体系(job/process.go); - Go 函数:
ExecutionTypeFunc在 Process 模式下不受支持,源码中会回退到 goroutine 执行器(job/worker.go)。
值得注意的细节是convertArgsForYaoRun(job/process.go):基础类型(字符串、布尔、整数、浮点)直接转为命令行参数,而 slice、map、struct 等复杂类型会被 JSON 序列化并以::前缀传入(如::{"a":1}),以便yao run正确解析结构化参数。
选型建议:短小、高频、对延迟敏感的任务优先选择GOROUTINE模式(共享内存、无进程启动开销);需要强隔离(内存泄漏隔离、崩溃不相互影响)或必须借助yao runCLI 能力的任务选择PROCESS模式。
四、三种调度类型与创建 API
调度类型由ScheduleType常量定义(job/types.go):once、cron、daemon。对应三个创建函数(job/job.go):
| 函数 | 调度类型 | 说明 |
|---|---|---|
Once(mode ModeType, data map[string]interface{}) (*Job, error) | once | 一次性任务,执行完毕即结束 |
Cron(mode ModeType, data map[string]interface{}, expression string) (*Job, error) | cron | 按 Cron 表达式周期执行,如"0 2 * * *"(每天凌晨 2 点) |
Daemon(mode ModeType, data map[string]interface{}) (*Job, error) | daemon | 持续运行的守护任务,配合ctx.Done()优雅退出 |
源码实现中,三个函数本质都是"设置mode与schedule_type字段后,将 data map 序列化为 JSON 并交给内部makeJob反序列化成Job结构体"(job/job.go)。makeJob还会做一系列默认值填充:Status默认draft、MaxWorkerNums默认 1、Priority默认 1、CreatedBy默认system、新任务默认Enabled = true。
此外还提供了三组"创建即保存"的快捷函数:OnceAndSave、CronAndSave、DaemonAndSave(job/job.go),它们创建任务后立即调用SaveJob落库。
4.1 快速上手:一次性任务
// 创建 Goroutine 模式的一次性任务 job, err := job.Once(job.GOROUTINE, map[string]interface{}{ "name": "Example Job", "description": "This is an example job", }) handler := func(ctx context.Context, execution *job.Execution) error { execution.Info("Job started") execution.SetProgress(50, "In progress...") // 执行业务逻辑 time.Sleep(1 * time.Second) execution.SetProgress(100, "Completed") execution.Info("Job completed") return nil } err = job.Add(1, handler) // 早期接口形式,见下方 API 演进说明 if err != nil { return err } // 启动任务 err = job.Start()API 演进说明:README 中的
Add(priority, handler)属于早期接口形态。当前仓库源码中,job/execution.go 已演进为按执行类型添加执行体的三个方法:Add(options *ExecutionOptions, processName string, args ...interface{})(执行 Yao 处理器)、AddCommand(options, command, args, env)(执行系统命令)、AddFunc(options, name, fn ExecutionFunc, args)(执行 Go 函数),其中ExecutionFunc签名是func(ctx *ExecutionContext) error(job/types.go),ExecutionContext携带Ctx(Go 上下文)、Execution(当前执行实例)与Args(函数参数)。
4.2 当前源码形式的完整示例
// 一次性任务:调用 Yao 处理器 opts := job.NewExecutionOptions().WithPriority(10).AddSharedData("sid", "session-001") err = j.Add(opts, "flows.example.cleanup", map[string]interface{}{"days": 30}) // 一次性任务:执行系统命令 err = j.AddCommand(nil, "rm", []string{"-rf", "/tmp/old"}, map[string]string{"TZ": "Asia/Shanghai"}) // 一次性任务:执行 Go 函数 err = j.AddFunc(nil, "myHandler", func(ctx *job.ExecutionContext) error { ctx.Execution.Info("Hello from func, args=%v", ctx.Args) return nil }, map[string]interface{}{"k": "v"}) // 推送执行 err = j.Push() // 等价于 README 中的 Start()其中ExecutionOptions支持链式调用(job/types.go):WithPriority(int)设置优先级(越高越先执行,见 job/job.go 的排序逻辑)、WithSharedData(map)/AddSharedData(key, value)注入共享数据。
4.3 定时任务与守护任务
// Cron 定时任务(每天凌晨 2 点执行),采用独立进程模式 cronJob, err := job.Cron(job.PROCESS, map[string]interface{}{ "name": "Cleanup Task", }, "0 2 * * *") err = cronJob.Add(nil, "flows.cleanup.run") err = cronJob.Push()// 守护任务:持续运行,通过 ctx.Done() 优雅退出 daemonJob, err := job.Daemon(job.GOROUTINE, map[string]interface{}{ "name": "Monitor Daemon", }) err = daemonJob.AddFunc(nil, "monitor", func(ctx *job.ExecutionContext) error { ticker := time.NewTicker(5 * time.Second) defer ticker.Stop() for { select { case <-ctx.Ctx.Done(): return ctx.Ctx.Err() case <-ticker.C: ctx.Execution.Info("Performing monitor check") } } }, nil) err = daemonJob.Push()4.4 任务生命周期方法
Job 提供的生命周期方法一览:
Push()(job/job.go):从数据库加载该任务的全部执行实例,按优先级降序排序,将任务状态置为ready并落库,然后逐一向全局 WorkerManager 提交执行;提交失败会汇总错误并清理上下文。多个执行实例从同一个任务上下文context.WithCancel派生,保证整体可取消;Stop()(job/job.go):将任务置为disabled,取消任务级上下文并逐个取消运行中的执行,将queued/running状态改写成cancelled并写入取消日志;Destroy()(job/job.go):先Stop()再置为deleted,彻底清理资源;- 链式配置方法:
SetName、SetDescription、SetCategory、SetMaxWorkerNums、SetMaxRetryCount、SetDefaultTimeout、SetConfig、SetStatus等(job/job.go),全部返回*Job便于链式调用。
五、实时进度追踪:三层联动机制
进度追踪贯穿"业务代码 → 执行器 → 数据库"三层,job/README.md 中描述的三大特性在源码中均有对应实现:
- 实时进度更新:业务侧通过
execution.SetProgress(progress, message)(job/execution.go)更新进度,该方法同时更新内存中的Execution.Progress、落库保存执行实例,并自动写一条Progress: 50% - In progress...级别的日志; - 数据库持久化:
Progress管理器(job/progress.go)内部以sync.RWMutex保护进度状态,Set()在更新内存值的同时,若已设置ExecutionID则加载执行实例并同步保存到数据库;Get()供查询当前进度; - 回调支持:Goroutine 模式执行 Yao 处理器时,
proc.WithCallback注册的回调会识别type == "progress"的数据包,通过extractProgressData(job/goroutine.go)解析progress与message字段,随后更新执行进度、写入日志并落库(job/goroutine.go)。
对于系统命令类执行,还提供了外部上报入口UpdateExecutionProgress(executionID, progressData)(job/goroutine.go):命令脚本可以通过 HTTP API 携带执行 ID 上报进度,框架会加载执行实例、解析progress/message并保存。
ProgressManager接口(job/interfaces.go)只定义了一个方法:Set(progress int, message string) error,Progress结构体即为其标准实现;任务对象通过job.Progress()(job/progress.go)获取进度管理器实例。
六、结构化日志系统:七级日志与双通道输出
LogLevel定义了完整的七级日志(job/types.go),按严重程度从高到低为:Panic、Fatal、Error、Warn、Info、Debug、Trace。
Execution提供一组便捷方法(job/execution.go):Info、Debug、Warn、Error、Fatal、Panic、Trace,它们最终汇聚到统一的Log(level, format, args...)(job/execution.go)。该方法的处理流程体现了"结构化 + 双通道"的设计:
- 将
LogLevel映射为字符串级别(Fatal/Panic映射为fatal,Trace映射为debug); - 构造包含
JobID、ExecutionID、WorkerID、ProcessID、当前Progress、Timestamp的Log记录,调用SaveLog写入数据库; - 同时输出到系统 logger,带上前缀
[Job:<id>][Exec:<id>],便于在服务端日志中快速检索定位。
执行过程中框架还会自动落库关键事件日志:执行开始/完成时写入info级别日志(含进度与耗时),失败时写入error级别日志(含 Worker ID 与错误信息),见 job/worker.go 的processWork逻辑。
七、Worker 管理系统:池化、分发与背压
7.1 全局单例与默认规模
WorkerManager 以全局单例形式存在(GetWorkerManager,job/worker.go),并在包初始化时自动Start()。默认 Worker 数量为runtime.NumCPU() * 4(job/worker.go),兼顾资源利用与系统负载;测试场景可通过NewWorkerManagerForTest(maxWorkers)创建非单例实例。
7.2 经典 Worker Pool 架构
源码实现了 Go 并发编程中最经典的 Worker Pool 模式(job/worker.go):
workQueue chan *WorkRequest:容量为maxWorkers*4的待处理队列,允许约 4 倍超额积压;workerPool chan chan *WorkRequest:空闲 Worker 的注册通道;activeWorkers map[string]*Worker:活跃 Worker 注册表;dispatch()(job/worker.go):调度协程从workQueue取出请求,再从workerPool取一个空闲 Worker 的 JobChannel 投递任务;Worker.Start()(job/worker.go):每个 Worker 循环执行"注册自己 → 等待任务 → 处理任务"。
7.3 提交与背压控制
SubmitJob(job/worker.go)在提交前检查队列水位:queueLen >= queueCap时直接返回"work queue is full (n/cap), please retry later"错误,实现背压保护,避免无界积压;随后以非阻塞 goroutine 写入队列,若任务上下文已取消则放弃提交。监控接口GetQueueStatus()返回当前队列长度与容量,GetActiveWorkers()返回活跃 Worker 数,可接入服务健康面板。
7.4 执行处理主流程
每个任务由processWork(job/worker.go)统一处理,流程为:置执行状态为running并记录worker_id、started_at→ 更新任务状态为running并记录current_execution_id、last_run_at→ 若配置了DefaultTimeout则用context.WithTimeout包裹执行上下文 → 按模式分发执行 → 计算duration与ended_at→ 成功则置completed并将进度置 100,失败则置failed并写入error_info→ 一次性任务完成时任务状态置completed,定时/守护任务则回到ready等待下次执行。
八、系统自愈能力:健康检查与数据清理
框架在包初始化时(job/job.go)启动两个后台守护组件:
- 健康检查器(job/health.go):默认每 5 分钟执行一次健康检查,通过
NewHealthChecker(interval)创建、StopHealthChecker()/RestartHealthChecker(interval)停止或动态调整间隔(job/job.go); - 数据清理器:默认保留 90 天执行/日志数据,
initDataCleaner创建NewDataCleaner(90)并启动;测试或运维场景可用ForceCleanup()强制执行清理(job/job.go)。
两者都支持通过测试或动态配置重启,为长时间运行的任务系统提供基础的自愈与容量管理能力。
九、数据库 CRUD API 速查
job/data.go 提供了一整套面向四个模型的分页查询、读写与删除函数,README 中列出的核心函数均可在源码中找到对应实现:
Job 相关
ListJobs(param model.QueryParam, page int, pagesize int) (maps.MapStrAny, error)(job/data.go):分页列出任务,并自动批量查询分类名填充category_name字段;GetJob(jobID string) (*Job, error)(job/data.go):按job_id查询并关联分类名;SaveJob(job *Job) error(job/data.go):新增(ID == 0,自动生成job_id)或更新(按job_id定位,保护id/created_at不被覆盖);若仅填了CategoryName会自动解析为分类 ID;RemoveJobs(ids []string) error、GetActiveJobs()、CountJobs(param)。
Category 相关
GetOrCreateCategory(name, description string) (*Category, error)(job/data.go):按名查询,不存在则自动创建;GetCategories(param)、SaveCategory、RemoveCategories、CountCategories。
Execution 相关
GetExecutions(jobID string) ([]*Execution, error)(job/data.go):按任务查执行实例,并从ConfigSnapshot还原执行配置;GetExecution(executionID, param)、SaveExecution、RemoveExecutions、CountExecutions。SaveExecution内部还做了 SQLite 兼容处理:将config_snapshot、execution_options、result、error_info等 JSON 字段显式转为字符串存储(job/data.go)。
Log 相关
ListLogs(jobID string, param model.QueryParam, page, pagesize)(job/data.go):按任务分页查日志,默认按timestamp倒序;SaveLog(log *Log) error(始终追加新记录)、RemoveLogs(ids []string) error。
十、测试覆盖与运行方式
框架遵循"每个功能模块都有对应测试文件"的工程规范(job/README.md 的 File Structure 与 Test Coverage 章节),实际文件与 job/ 目录一致:
- job/data_test.go + job/data_internal_test.go:Job / Category / Execution / Log 四类 CRUD 测试;
- job/job_test.go:任务管理集成测试;
- job/worker_test.go:Worker 生命周期、任务提交、双模式、错误处理、并发测试;
- job/progress_test.go:进度管理器、执行中进度、数据库持久化、进度查询测试;
- job/types_test.go:类型常量与结构体测试;
- job/health_test.go:健康检查器测试。
运行测试(job/README.md 的 Running Tests 章节):
# 运行全部 CRUD 测试 go test -v ./job/... -run "CRUD" # 运行全部类型测试 go test -v ./job/... -run "Types|Structure" # 运行 Worker 管理测试 go test -v ./job/... -run "Worker" # 运行进度管理测试 go test -v ./job/... -run "Progress" # 运行全部测试 go test -v ./job/...环境要求
测试依赖数据库模型与平台配置,运行前需先加载环境变量(job/README.md 的 Environment Requirements 章节):
source $YAO_ROOT/env.local.sh由于框架持久化依赖__yao.job系列内置模型,请确保应用已通过yao migrate完成模型建表(相关模型见 yao/models/job/)。
十一、架构总结与选型要点
综合 job/README.md 与源码,Job Framework 的 8 项技术特征可归纳为:
- 完整 CRUD:四个模型全量增删改查,均有测试覆盖;
- 双执行模式:Goroutine 轻量快速、独立进程强隔离,模式可随任务指定;
- 自动分类管理:按名 Get-or-Create,默认分类开箱即用;
- 实时进度追踪:业务代码、Yao 回调、外部 API 三条通道均可更新进度;
- 结构化日志:七级日志、执行上下文完备、数据库 + 系统双通道输出;
- 并发安全:Worker Pool + 队列背压 + 互斥锁保护,支持多 Worker 并发执行;
- 数据持久化:任务、执行实例、进度、日志全部落库,重启可恢复;
- 完备测试:每个模块均有对应单元/集成测试文件。
实战选型建议:
- 高频轻量任务(心跳监控、缓存刷新)→
GOROUTINE+Once/Daemon; - 强隔离或需 CLI 能力(数据清理、批处理脚本)→
PROCESS+Cron; - 需要跨执行传递用户会话时,使用
ExecutionOptions.WithSharedData注入sid、authorized等上下文(Goroutine 模式会还原为进程的 Global/SID/Authorized 环境,见 job/goroutine.go); - 需要外部脚本上报进度时,通过
UpdateExecutionProgress提供执行 ID 即可接入; - 高并发场景留意
SubmitJob的队列满错误,可结合GetQueueStatus()监控积压水位并采取退避重试策略。
- Agent 框架
- 后端
- 低代码
- RAG
【免费下载链接】yao
✨ All your agents and workspaces in one place, on every device you own. Track tasks on a board, accessible from desktop, mobile, browser, or API. Self-hosted.
相关推荐
Delta 金手指模拟器 3 种用法:预设库、手输代码与不生效排查
Delta 金手指模拟器 3 种用法:预设库、手输代码与不生效排查 Delta 是一款跑在 iOS 上的多平台经典游戏模拟器,它的金手指(Cheat Code)
游戏开发PowerJob vs XXL-Job:新一代分布式任务调度框架深度测评
PowerJob vs XXL Job:新一代分布式任务调度框架深度测评 在分布式系统架构中,任务调度框架扮演着关键角色,负责协调各类定时任务、批量处理和分布式
任务调度后端Celery 任务执行追踪机制深度解析:celery.app.trace 模块源码级指南
Celery 任务执行追踪机制深度解析:celery.app.trace 模块源码级指南 导读 本文聚焦 Celery 分布式任务队列中承载任务执行追踪(tas
任务调度后端消息队列
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考