Skip to content

DAG 与耐久编排 ​

问题:进程重启后,长任务从哪一步继续 ​

一个导出过程等待文件准备、计算、审核、上传。单进程协程可以运行它,但重启会丢失等待位置;任务跨数小时后,简单的 HTTP 超时更无法代表业务时限。

耐久编排将推进过程记录为可恢复历史。DAG 则描述步骤之间的依赖:只有依赖满足,节点才能运行。它们可以配合,但有 DAG 不等于耐久,有持久队列也不等于完整编排。

原理:图、历史与事实各有职责 ​

DAG 必须拒绝重复步骤、未知依赖与环,否则任务可能永远等待。输入映射决定上一步输出怎样进入下一步,应该校验并使用发布时冻结的定义。

Temporal 中 Workflow 描述耐久控制流程;Activity 执行网络与数据库等外部协作;Signal 将外部结果送回等待中的流程。Workflow 会根据历史重放,因此使用引擎提供的确定性时间、定时器和并发原语,避免直接读当前时间或随意发网络请求。

数据库仍可拥有业务查询与状态规则,Workflow 历史拥有编排恢复依据。两者应通过幂等 Activity 协作,不能让两个事实源各自无条件决定最终状态。

最小例子:编排意图 ​

text
教学伪代码:
source = Activity("download", stableStepID)
result = Activity("transform", source, stableStepID)
Activity("upload", result, stableStepID)
等待审核 Signal 或耐久 Timer
Activity("finalize", stableRunID)

Activity 成功但确认历史写入前失败可能重试。stableStepID 不能每次随机生成;下游创建 Attempt 或上传结果要能识别同一次动作。Workflow 重放与 Activity 重试是两种机制,不能混为“重新执行所有代码”。

DAG:依赖表达可并行的部分
100%
下载源文件
→
校验内容
生成缩略图
并行
提取元数据
两个分支都完成
→
发布结果

可触控平移或键盘滚动;每张图的正文同时提供文字解释。

文字解释:校验完成后可以并行处理两个分支,发布依赖两个结果。图中的合流必须定义失败、跳过和取消语义,不能只画成功路径。

方案比较 ​

实现适用场景需要自行管理
内存 goroutine短任务、允许丢失重启与超时状态
数据库队列+状态机较简单持久步骤定时器、依赖、恢复扫描
耐久工作流引擎长流程、等待、多步重试Activity 幂等、版本与运维

引擎不会解决全部业务语义:失败后回滚外部动作仍需要补偿,已经发出的邮件不能由 Workflow rollback 自动收回。

真实案例:DAG 校验与执行快照 ​

当前 DAG 校验用三色访问状态检测环,先检查已知步骤与依赖。依赖判定允许 SUCCEEDED 或 SKIPPED,表明“跳过”在此实现中是可满足依赖的控制结果,需要结合条件与输出映射理解。

已核对的实现 · 本地代码快照

DAG 依赖与环检测 backend · c6dbe05c

server/internal/adapters/inbound/temporal/dag_workflow.go · 第 45–111 行
符号:validateStepPlans · 核对日期 2026-10-02
来源与提交版本一致

func validateStepPlans(steps []StepPlan) error {
	known := make(map[string]StepPlan, len(steps))
	for _, step := range steps {
		if step.Key == "" || step.StepID == "" {
			return errors.New("step key and ID are required")
		}
		if _, exists := known[step.Key]; exists {
			return fmt.Errorf("duplicate step key %q", step.Key)
		}
		known[step.Key] = step
	}
	for _, step := range steps {
		for _, dependency := range step.DependsOn {
			if _, exists := known[dependency]; !exists {
				return fmt.Errorf("step %q depends on unknown step %q", step.Key, dependency)
			}
		}
	}
	state := make(map[string]uint8, len(steps))
	var visit func(string) error
	visit = func(key string) error {
		switch state[key] {
		case 1:
			return fmt.Errorf("dependency cycle includes step %q", key)
		case 2:
			return nil
		}
		state[key] = 1
		for _, dependency := range known[key].DependsOn {
			if err := visit(dependency); err != nil {
				return err
			}
		}
		state[key] = 2
		return nil
	}
	for key := range known {
		if err := visit(key); err != nil {
			return err
		}
	}
	return nil
}

func dependenciesComplete(step StepPlan, states []dagStepState, byKey map[string]int) bool {
	for _, dependency := range step.DependsOn {
		status := states[byKey[dependency]].status
		if status != dagStepSucceeded && status != dagStepSkipped {
			return false
		}
	}
	return true
}

func hasRunningSteps(states []dagStepState) bool {
	for _, state := range states {
		if state.status == dagStepRunning {
			return true
		}
	}
	return false
}

func firstRunningStepID(steps []StepPlan, states []dagStepState) string {
	for index, state := range states {
		if state.status == dagStepRunning {
			return steps[index].StepID

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

实体、值对象与创建不变量 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
}

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

Execution 保存 workflow_snapshot,避免运行中定义修改影响本次执行。Worker 上报先进入 Server Inbox,再推动状态与 Signal;NATS 起桥接作用,不替代数据库事实。

失败与边界:流程升级与回放 ​

修改已运行 Workflow 的控制分支可能让历史重放不一致。升级要采用引擎支持的版本策略并验证 Replay;仅通过普通单元测试不足以证明历史兼容。此站不会自动执行真实工作流。

取消要传递到运行中的子步骤和 Worker,还要处理发送失败及迟到事件。定时器竞争成功事件时,终态应通过数据库规则收敛。

迁移练习与参考答案 ​

练习:证件审核下载图片、并行 OCR 与风控、合并结果后等待人工审批。如何避免重启丢失等待?

参考答案:冻结流程定义,使用稳定步骤身份与持久编排;OCR 和风控为 Activity 或外部 Step,合流明确两个结果的失败策略,人工审批用带审核身份的 Signal,超时用耐久 Timer。审批来源需要认证,重复 Signal 不能重复发布结果。

继续阅读:Worker 恢复、测试与故障验证。

理解原理 · 分析取舍 · 用真实代码检验