切换主题
任务模型与状态机
问题:失败后把状态改回待执行就够了吗
报表导出失败后,用户点击重新执行。如果只把同一行从 FAILED 改成 PENDING,上次用了哪版模板、在哪一步失败、是否已生成部分文件都会被覆盖。迟到的旧结果还可能被当成新运行成功。
长期意图和一次运行应有不同身份。进一步把运行拆成步骤,把每次实际执行拆成尝试,可以分别回答“用户要什么”“本次怎样做”“哪个动作”“第几次做”。
原理:状态是对事实的压缩表达
状态机明确哪些转换被允许,转换需要什么证据,以及终态是否可以再被普通事件改变。不要用一组可随意组合的 running/failed/canceled 布尔值,让矛盾状态变得可表示。
取消应区分“已提出取消”与“确认停止”。已经开始的外部动作无法总是立即撤销。CANCELING 是中间事实,CANCELED 需要确认或明确的失效判定。
最小例子:状态转换是一个领域行为
go
// 教学片段:持久化版本检查在 repository 中完成。
func (run *Run) Finish(result Result) error {
if run.State != "RUNNING" { return errors.New("run is not running") }
run.State = "SUCCEEDED"
run.Result = result
run.Version++
return nil
}状态方法保护单对象规则;保存时必须验证读到的版本仍有效。同一次成功事件重投应返回同一结果;失败或取消后收到不同终态要遵循明确竞争策略。
交互 04
取消是请求,终态是事实
Execution #1 · CREATED
- Execution #1 · CREATED
转换集合来自 Execution 领域模型。重跑入口在本教学模型中限制为终态,具体权限与业务条件仍由应用用例决定。这里没有模拟 Worker 停止、定时器或外部副作用;CANCELED 的落库必须有相应证据。
尝试先进入 CANCELING,观察为什么不能再直接 SUCCEEDED;再创建新执行,观察上次状态历史如何保留。模型允许某些早期状态直接终结,这是当前案例的转换规则,其他业务可以不同。
方案比较
单一 job 行适合简单、短生命周期工作,但重试审计和步骤恢复困难。Task/Run/Step/Attempt 分层提供历史与幂等边界,也增加查询和状态聚合成本。不要为了统一命名强迫每个通知任务都建立四层。
状态机可由代码枚举、数据驱动定义或耐久工作流引擎表达。数据库约束、领域方法与编排器必须保持一致;图中有一条箭头不代表任意时刻都可以通过 API 强制转换。
真实案例:四层执行模型
TDP 的 Task 保存长期请求及 CurrentExecutionID;Execution 保存运行序号和 workflow_snapshot;Step 表达 DAG 节点;Attempt 对应一次 Worker 执行。内部重试创建新 Attempt,用户重跑产生新的 Execution。
已核对的实现 · 本地代码快照
实体、值对象与创建不变量 backend · c6dbe05c
server/internal/core/execution/domain/model.go · 第 12–141 行
符号:NewTask · 核对日期 2026-10-02
来源与提交版本一致
// ID 是 Execution 上下文使用的强类型资源标识,避免把普通字符串误传为领域 ID。
type ID string
// ParseID 去除首尾空白并拒绝空标识。
func ParseID(value string) (ID, error) {
value = strings.TrimSpace(value)
if value == "" {
return "", errors.New("ID is required")
}
return ID(value), nil
}
// Document 保存已经校验过的 JSON 文档,并在进入领域层时复制以避免外部修改。
type Document []byte
// ParseDocument 校验并复制 JSON 字节,调用方可安全复用原缓冲区。
func ParseDocument(value []byte) (Document, error) {
if len(value) == 0 || !json.Valid(value) {
return nil, errors.New("document must contain valid JSON")
}
return append(Document(nil), value...), nil
}
// Task 是业务请求的长期聚合根;一次 Task 可以通过重跑产生多个 Execution。
type Task struct {
ID ID // 平台生成的任务标识。
ApplicationID ID // 任务的归属和访问范围。
Name string // 供业务人员识别的任务名称。
ExternalRef string // 业务系统用于关联和幂等追踪的外部引用。
WorkflowDefinitionID ID // 创建任务时解析到的已发布流程。
Input Document // 业务系统提交给流程的原始输入。
StepConstraints Document // 按步骤合并后的调度约束。
EffectiveStatus ExecutionStatus // 当前执行状态,便于直接查询任务。
CurrentExecutionID ID // 最近一次创建或重跑的执行。
Priority int16 // 待调度任务的业务优先级,范围为 -100 到 100。
Version int64 // 用于任务的乐观并发控制。
CreatedAt time.Time // 任务首次创建时间。
UpdatedAt time.Time // 聚合最后一次状态变更时间。
}
// Execution 是 Task 的一次具体运行,保存创建时冻结的流程快照。
type Execution struct {
ReconciliationError string // 终态对账失败或待人工检查的原因。
SchedulingReasons []string // 当前排队步骤的等待原因。
ID ID // 本次执行的稳定标识。
TaskID ID // 所属任务。
Number int32 // 同一任务下从 1 开始递增的执行序号。
Status ExecutionStatus // 执行生命周期状态。
Progress int16 // 0 到 100 的单调进度值。
Result Document // 执行完成后的 JSON 结果。
WorkflowSnapshot Document // 冻结本次运行使用的已编译流程。
CallbackURL string // 完成事件的业务回调地址。
Reason string // 记录重跑、取消或失败等人工可读原因。
Version int64 // 用于执行的乐观并发控制。
CreatedAt time.Time // 执行创建时间。
StartedAt *time.Time // 首次进入运行态的时间,未开始时为空。
CompletedAt *time.Time // 进入终态的时间,未结束时为空。
}
// NewTaskParams 汇集原子创建 Task 及首个 Execution 所需的已解析参数。
type NewTaskParams struct {
TaskID ID // 预先生成的任务标识。
ExecutionID ID // 首个执行标识。
ApplicationID ID // 任务所属应用。
Name string // 任务展示名称。
ExternalRef string // 业务系统外部引用。
WorkflowDefinitionID ID // 已发布流程标识。
Input Document // 已经校验的任务输入。
WorkflowSnapshot Document // 创建时冻结的流程定义。
StepConstraints Document // 已经合并的步骤调度约束。
CallbackURL string // 可选的执行事件回调地址。
Priority int16 // 任务优先级。
CreatedAt time.Time // 统一注入的创建时间。
}
// NewTask 校验创建不变量,并同时构造 Task 与编号为 1 的 Execution。
func NewTask(params NewTaskParams) (*Task, *Execution, error) {
if err := validateIDs(params.TaskID, params.ExecutionID, params.ApplicationID, params.WorkflowDefinitionID); err != nil {
return nil, nil, err
}
if strings.TrimSpace(params.Name) == "" {
return nil, nil, errors.New("task name is required")
}
if len([]rune(strings.TrimSpace(params.Name))) > 200 {
return nil, nil, errors.New("task name must not exceed 200 characters")
}
if strings.TrimSpace(params.ExternalRef) == "" {
return nil, nil, errors.New("external reference is required")
}
if params.Priority < -100 || params.Priority > 100 {
return nil, nil, errors.New("priority must be between -100 and 100")
}
if len(params.StepConstraints) == 0 {
params.StepConstraints = Document(`{}`)
}
if !json.Valid(params.Input) || !json.Valid(params.WorkflowSnapshot) || !json.Valid(params.StepConstraints) {
return nil, nil, errors.New("input and workflow snapshot must contain valid JSON")
}
if params.CreatedAt.IsZero() {
return nil, nil, errors.New("creation time is required")
}
execution := &Execution{
ID: params.ExecutionID,
TaskID: params.TaskID,
Number: 1,
Status: ExecutionCreated,
WorkflowSnapshot: append(Document(nil), params.WorkflowSnapshot...),
CallbackURL: params.CallbackURL,
Version: 1,
CreatedAt: params.CreatedAt,
}
task := &Task{
ID: params.TaskID,
ApplicationID: params.ApplicationID,
Name: strings.TrimSpace(params.Name),
ExternalRef: strings.TrimSpace(params.ExternalRef),
WorkflowDefinitionID: params.WorkflowDefinitionID,
Input: append(Document(nil), params.Input...),
StepConstraints: append(Document(nil), params.StepConstraints...),
EffectiveStatus: execution.Status,
CurrentExecutionID: execution.ID,
Priority: params.Priority,
Version: 1,
CreatedAt: params.CreatedAt,
UpdatedAt: params.CreatedAt,
}
return task, execution, nil
}
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
执行状态机及单调进度 backend · c6dbe05c
server/internal/core/execution/domain/status.go · 第 22–79 行
符号:Transition / SetProgress · 核对日期 2026-10-02
来源与提交版本一致
var executionTransitions = map[ExecutionStatus]map[ExecutionStatus]struct{}{
ExecutionCreated: set(ExecutionQueued, ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceling, ExecutionCanceled),
ExecutionQueued: set(ExecutionRunning, ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceling, ExecutionCanceled),
ExecutionRunning: set(ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceling, ExecutionCanceled),
ExecutionCanceling: set(ExecutionCanceled, ExecutionFailed, ExecutionTimedOut),
}
// Terminal 判断执行是否已经结束且不会再接受普通状态推进。
func (s ExecutionStatus) Terminal() bool {
switch s {
case ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceled:
return true
default:
return false
}
}
// Transition 按状态机推进执行,并统一维护开始时间、完成时间和版本号。
func (e *Execution) Transition(to ExecutionStatus, at time.Time) error {
if e.Status == to {
return nil
}
if _, ok := executionTransitions[e.Status][to]; !ok {
return fmt.Errorf("execution transition %s -> %s is not allowed", e.Status, to)
}
if to == ExecutionRunning && e.StartedAt == nil {
startedAt := at
e.StartedAt = &startedAt
}
if to.Terminal() {
completedAt := at
e.CompletedAt = &completedAt
if to == ExecutionSucceeded {
e.Progress = 100
}
}
e.Status = to
e.Version++
return nil
}
// SetProgress 仅允许运行中的执行单调增加 0 到 100 的进度。
func (e *Execution) SetProgress(progress int16) error {
if e.Status != ExecutionRunning {
return fmt.Errorf("cannot update progress while execution is %s", e.Status)
}
if progress < e.Progress || progress > 100 {
return fmt.Errorf("progress must be monotonic and between %d and 100", e.Progress)
}
e.Progress = progress
e.Version++
return nil
}
// StepStatus 表示工作流步骤从等待依赖到完成的生命周期状态。
type StepStatus string
const (片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
Execution.Transition 对同状态重复更新直接返回,不重复递增版本;进入 RUNNING 设置首次开始时间,进入终态设置完成时间。SetProgress 只在 RUNNING 允许单调增加,成功终态设置为 100。
失败与边界:进度不等于业务成功
进度可以是估算值,100% 的文件处理仍可能在保存结果时失败。成功应以持久化终态及必要输出为依据。多个并发步骤的总体进度需明确权重,避免重试导致进度倒退或重复累计。
超时说明协调者不再等待,并不证明远端进程已结束。Attempt 租约、事件来源版本和资源回收需要共同管理。终态不可回退,也不意味着不再做对账或清理。
迁移练习与参考答案
练习:视频转码有下载、编码、上传三步;编码失败重试,用户又发起整次重跑。如何分配标识?
参考答案:视频业务请求身份稳定;每次整次运行一个 Run;三步分别有 Step;编码每次重试有不同 Attempt。产物路径或远端请求应含相应执行身份,旧 Attempt 的迟到事件不能覆盖新 Attempt。保留旧 Run 失败记录供诊断。