2026-09-19-task-orchestration-design.md 18 KB

任务编排系统设计

优先级: P1 预计工时: 后端 5d + 前端 3d 状态: 设计已确认(待写实现计划) 日期: 2026-09-19

1. 背景与问题

现有任务系统(tasks 表)已支持单前置(prerequisite_task_id)、循环(repeat_type)、父链(parent_task_id),但缺少:

  • 多前置 AND/OR:「数学+英语+阅读三项全部完成后,才解锁奖励任务」
  • 多前置 OR:「数学或英语任一完成后,即可开始下一项」
  • 后置兜底:「若晨读在 30 分钟内未完成,自动派发补救任务」
  • 带终止条件的循环:「晨读最多重试 3 次,3 次仍未完成则标记失败并触发后续计划」

这些能力需要新增上层编排模块,不修改现有 tasks 表结构,复用现有 TaskService 的创建/进度/积分体系。


2. 核心设计决策(已与用户确认)

决策项 结论
驱动场景 家庭/成长任务的条件编排
最大复杂度 含环 DAG
前置语义 AND(全部) + OR(任一)均支持
未完成语义 超时未按时完成 + 主动标记失败,均支持
流定义方式 UI 可视化编排(jsPlumb 画布)
节点与 tasks 关系 节点执行时创建 tasks 实例,复用现有进度/积分体系
循环终止 计数终止 + 条件终止两者都支持
调度模型 引入 Quartz 独立轮询,不依赖 @Scheduled
内部回调范围 仅编排节点任务回调,普通任务零开销
画布选型 jsPlumb(Element UI 兼容,vue2 有现成封装)
版本管理 执行实例快照旧版本,新执行用新版本
并发保护 乐观锁 + 幂等检查

3. 数据模型

3.1 task_orchestration_flows(编排流定义)

字段 类型 说明
id BIGINT PK
name VARCHAR(100) 流名称,如「晨读打卡补救流」
description TEXT 描述
creator_id BIGINT 创建者(规划师/家长)
family_id BIGINT 所属家庭(null = 系统级公共模板)
version INT 版本号,每次 publish 递增
status VARCHAR(20) draft / published / archived
schedule_cron VARCHAR(64) 定时触发 cron(可为空,仅手动触发时为空)
config_json JSON 流的核心配置(节点列表 + 边的初始快照,用于审计和回溯)
created_at / updated_at DATETIME

3.2 task_orchestration_edges(依赖边)

字段 类型 说明
id BIGINT PK
flow_id BIGINT 所属流(FK → task_orchestration_flows)
from_node_id VARCHAR(64) 源节点 ID(与 config_json.nodes[].id 对应)
to_node_id VARCHAR(64) 目标节点 ID
edge_type VARCHAR(20) success / timeout / failed
operator VARCHAR(10) AND / OR
sort_order INT 同 to_node 多条边的排序权重
created_at DATETIME

3.3 task_orchestration_executions(执行实例)

字段 类型 说明
id BIGINT PK
flow_id BIGINT 所属流(FK)
flow_version INT 触发时对应的流版本号(用于版本快照隔离)
family_id BIGINT 所属家庭
family_member_id BIGINT 绑定的家庭成员
trigger_source VARCHAR(20) manual / scheduled / api
status VARCHAR(20) running / paused / completed / failed / terminated
started_at / finished_at DATETIME
error_reason VARCHAR(255) 终止原因(terminated 时记录)
created_at / updated_at DATETIME

3.4 task_orchestration_node_instances(节点执行实例)

字段 类型 说明
id BIGINT PK
execution_id BIGINT 所属执行实例(FK)
node_id VARCHAR(64) 对应 flow 中的节点 ID
task_id BIGINT 关联现有 tasks.id(创建后写入,pending 时为 null)
status VARCHAR(20) pending / started / in_progress / completed / failed / timeout / skipped / terminated
generation INT 循环次数(首次 = 0,重试 = 1, 2...)
condition_met_at DATETIME 条件终止:外部 API 通知条件达成的时间(null = 未达成)
loop_termination_reason VARCHAR(50) max_loops_reached / condition_met / 空(未终止)
started_at / completed_at DATETIME
error_reason VARCHAR(255) 失败原因(可选)
created_at / updated_at DATETIME

唯一索引:(execution_id, node_id, generation) — 防止重复创建。


4. 节点 config_json 结构

每个节点的完整配置(存储在 flows.config_json.nodes[].config):

{
  "id": "node_1",
  "title": "晨读打卡",
  "description": "每日晨读 15 分钟",
  "task_template_ref": "tmpl_123",
  "is_start_node": true,
  "timeout_minutes": 30,
  "max_loops": 3,
  "loop_termination": {
    "type": "count",   // count | condition
    "condition_flag": null  // type=condition 时指定外部检查标识(在 node_instance 的 condition_met_at 字段中置值),由 API `/api/orchestration/node/condition-met` 置值
  },
  "retry_on_failure": false
}

5. 执行引擎核心机制

5.1 节点触发规则(每次评估时重新计算)

对每个 pending 节点,检查其所有入边:

operator 触发条件 示例
AND 所有入边源节点均为 completed A AND B → C(A、B 都完成才触发 C)
OR 任一入边源节点为 completed A OR B → C(A 或 B 任一完成就触发 C)

无入边的节点(is_start_node = true)在 flow 启动时直接触发。

5.2 循环处理

  • 节点被重新触发时,generation++
  • 计数终止:若 generation >= max_loops,节点状态置为 terminated,loop_termination_reason = max_loops_reached,触发所有 edge_type=timeout 的出边(兜底)
  • 条件终止:节点配置 loop_termination.type = condition 时,引擎每次轮询检查当前 generation 节点实例的 condition_met_at 是否已置值:
    • 已置值 → 节点置 completed,loop_termination_reason = condition_met,触发所有 edge_type=success 的出边
    • 未置值且 generation >= max_loops → 置 terminated,loop_termination_reason = max_loops_reached,触发 timeout 出边
  • 环中其他节点不受影响,继续正常评估

5.3 后置兜底触发

触发来源 节点状态变化 边类型 处理
deadline 到期且未完成 pending → timeout edge_type=timeout 创建目标节点任务实例
用户手动标记放弃 pending → failed edge_type=failed 创建目标节点任务实例
异常(如任务模板不存在) pending → failed edge_type=failed 同上

节点标记为 skipped 时(flow 终止/暂停),不触发任何出边。

5.4 节点创建任务逻辑

当节点触发时:
  1. 从 task_template_ref 加载任务模板配置
  2. 调用 TaskService.createTask() 创建 tasks 记录
  3. 更新 node_instance.task_id = tasks.id
  4. 设置 node_instance.status = "in_progress"
  5. deadline = now() + timeout_minutes

5.5 内部回调接口(供 TaskService 调用)

三个方法仅对属于编排节点的任务生效(先查 task_orchestration_node_instances.task_id = taskId,找不到则直接返回,零开销):

方法 时机 效果
onTaskCompleted(Long taskId) 任务完成/审核通过 节点 → completed,评估下游
onTaskFailed(Long taskId) 用户主动放弃/异常 节点 → failed,触发 failed 边
onTaskTimeout(Long taskId) deadline 到期未完成 节点 → timeout,触发 timeout 边

5.6 并发安全

  • evaluateFlow(executionId) 对同一个 execution_id 加分布式锁(Redis 或数据库行锁),避免 Quartz 轮询与回调同时触发导致重复创建节点任务
  • 节点实例唯一索引 (execution_id, node_id, generation) 作为最终防线

6. 调度与轮询

6.1 Quartz Job:OrchestrationPollingJob

@DisallowConcurrentExecution
public class OrchestrationPollingJob implements Job {
    @Override
    public void execute(JobExecutionContext context) {
        // 1. 查询所有 running 状态的 execution
        // 2. 批量检查 timeout 节点(deadline 到期)
        // 3. 批量调用 evaluateFlow 评估 pending 节点
        // 4. 检查 condition 终止节点的条件是否达成
        // 5. 更新 execution 状态(若所有节点 completed/skipped 则置 completed)
    }
}
  • 触发频率:每 30 秒执行一次(通过 Quartz TriggerBuilder 配置)
  • @DisallowConcurrentExecution:同一 job 实例不并发执行
  • 执行窗口:每次执行最多处理 N 个 execution(默认 20),避免单次执行过长

6.2 事件驱动(补充轮询)

除 Quartz 轮询外,以下节点状态变更即时触发评估:

  • onTaskCompleted / onTaskFailed 回调(事件驱动,无需等待下一个轮询周期)
  • 手动 node/fail 接口调用

7. API 设计

统一 @PostMapping,路由前缀 /api/orchestration。

7.1 Flow 管理

接口 说明 请求体关键字段
/api/orchestration/flow/save 保存/更新流(草稿) name, description, family_id, config_json, schedule_cron
/api/orchestration/flow/publish draft → published,version+1 flow_id
/api/orchestration/flow/list 分页列表 page, pageSize, familyId, status
/api/orchestration/flow/detail 流详情(含 edges) flow_id
/api/orchestration/flow/archive published → archived flow_id
/api/orchestration/flow/delete 删除 draft 流 flow_id

7.2 执行管理

接口 说明 请求体关键字段
/api/orchestration/execution/start 手动启动流 flow_id, family_member_id
/api/orchestration/execution/pause 暂停执行 execution_id
/api/orchestration/execution/resume 恢复执行 execution_id
/api/orchestration/execution/terminate 终止整个流 execution_id, reason
/api/orchestration/execution/detail 执行详情(节点树) execution_id
/api/orchestration/execution/list 执行历史列表 family_member_id, page, pageSize

7.3 节点管理

接口 说明 请求体关键字段
/api/orchestration/node/fail 手动标记节点为失败 node_instance_id, reason
/api/orchestration/node/restart 重启节点(生成新任务实例,generation 从 0 重置) node_instance_id, reason
/api/orchestration/node/condition-met 通知条件达成(type=condition 循环节点专用) node_instance_id

7.4 内部回调(非 HTTP,Service 直接调用)

见 §5.5。


8. 前端设计(cfc-web 管理端)

8.1 页面路由

  • /orchestration — 编排流管理首页(列表 + 新建入口)
  • /orchestration/edit/:id — 编排画布编辑器
  • /orchestration/execution/:id — 执行详情只读页

8.2 编排画布(FlowEditor.vue)

  • 左侧节点面板:可拖拽的任务节点模板列表(从 admin_task_templates 读取)
  • 中央画布:jsPlumb 实现,支持节点拖拽、连线、删除连线
  • 右侧属性面板:选中节点/边后编辑配置(超时时间、max_loops、终止条件、AND/OR 等)
  • 发布校验:
    • 无孤立节点
    • 每条出边都有目标节点
    • 有环的流必须有 max_loops ≥ 1 或 condition 终止配置
    • 节点引用的 task_template_ref 必须存在
    • 有且仅有 1 个 is_start_node = true

8.3 执行详情(ExecutionDetail.vue)

  • 树形展示所有节点实例,每节点显示:状态徽章 + 关联 tasks 记录 + 创建时间
  • 时间线视图:按时间轴展示节点完成顺序
  • 手动干预:标记失败、终止流

9. 错误处理与边界情况

场景 处理策略
节点引用的任务模板被删除 发布校验拦截;运行时节点标记 failed,触发 failed 边
环内节点无限重试 generation >= max_loops 强制终止,标记 terminated,触发 timeout 兜底
用户终止执行 所有 pending 节点标记 skipped;运行中 tasks 保留不删除,积分只发已完成的
flow 重新发布(版本升级) 运行中的执行继续使用旧版本快照(flow_version 字段隔离);新执行用新版本
家庭删除/成员离开 执行实例级联终止,清理 pending 节点
并发评估 Quartz job + Redis 锁,evaluateFlow 方法级锁
失败重试 节点失败不自动重试;手动调用 /api/orchestration/node/restart 接口(重置 generation=0,重新创建 task),或调用 /api/orchestration/node/condition-met 通知条件达成(type=condition 循环)

10. 测试策略

10.1 单元测试(JUnit)

  • OrchestrationEngine.evaluateFlow():覆盖 AND/OR 语义、无前置节点、多入边混合
  • 循环终止:count 终止(max_loops 精确值)、condition 终止(flag 提前/延迟达成)
  • 幂等:同时触发两个回调时不重复创建节点任务

10.2 集成测试

  • 完整流执行:start → node1 完成 → node2 触发 → node3 超时 → 兜底 node4 完成 → flow 结束
  • 版本隔离:execution_v1 运行中时 publish v2,v1 不受影响
  • 并发测试:模拟 2 个回调同时到达,验证只有 1 次任务创建

10.3 前端测试

  • 画布拖拽连线正确性(jsPlumb 事件绑定)
  • 发布校验拦截无效配置

11. 数据库迁移

在 DatabaseInitializer.runMigrations() 新增迁移:

CREATE TABLE IF NOT EXISTS `task_orchestration_flows` (
  `id` bigint NOT NULL AUTO_INCREMENT,
  `name` varchar(100) NOT NULL,
  `description` text,
  `creator_id` bigint NOT NULL,
  `family_id` bigint DEFAULT NULL,
  `version` int NOT NULL DEFAULT 1,
  `status` varchar(20) NOT NULL DEFAULT 'draft',
  `schedule_cron` varchar(64) DEFAULT NULL,
  `config_json` json DEFAULT NULL,
  `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
  `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  PRIMARY KEY (`id`),
  KEY `idx_family_id` (`family_id`),
  KEY `idx_status` (`status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE IF NOT EXISTS `task_orchestration_edges` (
  `id` bigint NOT NULL AUTO_INCREMENT,
  `flow_id` bigint NOT NULL,
  `from_node_id` varchar(64) NOT NULL,
  `to_node_id` varchar(64) NOT NULL,
  `edge_type` varchar(20) NOT NULL,
  `operator` varchar(10) NOT NULL DEFAULT 'AND',
  `sort_order` int NOT NULL DEFAULT 0,
  `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
  PRIMARY KEY (`id`),
  KEY `idx_flow_id` (`flow_id`),
  UNIQUE KEY `uk_flow_edge` (`flow_id`, `from_node_id`, `to_node_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE IF NOT EXISTS `task_orchestration_executions` (
  `id` bigint NOT NULL AUTO_INCREMENT,
  `flow_id` bigint NOT NULL,
  `flow_version` int NOT NULL DEFAULT 1,
  `family_id` bigint NOT NULL,
  `family_member_id` bigint NOT NULL,
  `trigger_source` varchar(20) NOT NULL DEFAULT 'manual',
  `status` varchar(20) NOT NULL DEFAULT 'running',
  `started_at` datetime NOT NULL,
  `finished_at` datetime DEFAULT NULL,
  `error_reason` varchar(255) DEFAULT NULL,
  `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
  `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  PRIMARY KEY (`id`),
  KEY `idx_execution_flow` (`flow_id`),
  KEY `idx_execution_status` (`status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE IF NOT EXISTS `task_orchestration_node_instances` (
  `id` bigint NOT NULL AUTO_INCREMENT,
  `execution_id` bigint NOT NULL,
  `node_id` varchar(64) NOT NULL,
  `task_id` bigint DEFAULT NULL,
  `status` varchar(20) NOT NULL DEFAULT 'pending',
  `generation` int NOT NULL DEFAULT 0,
  `condition_met_at` datetime DEFAULT NULL COMMENT '条件终止:外部API通知条件达成时间',
  `loop_termination_reason` varchar(50) DEFAULT NULL,
  `started_at` datetime DEFAULT NULL,
  `completed_at` datetime DEFAULT NULL,
  `error_reason` varchar(255) DEFAULT NULL,
  `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
  `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_node_instance` (`execution_id`, `node_id`, `generation`),
  KEY `idx_execution_id` (`execution_id`),
  KEY `idx_task_id` (`task_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

同步更新 schema.sql 建表语句。


12. 范围边界(明确不做)

不做 原因
家庭挑战/五维打卡合并进编排 已有设计文档明确不合并
前端小程序端编排编辑 编排编辑只放 cfc-web 管理端;小程序端只做执行结果展示(复用 tasks 列表)
节点级精确 cron 定时 只用 flow 级 schedule_cron,节点触发完全由依赖边决定
跨家庭共享执行实例 执行实例严格绑定 family,系统级 flow 仅用于多个家庭复用定义
流市场/模板市场 后续迭代
节点级并行执行(同时创建多个任务) 当前版本不支持,一个节点同一时刻只有一个任务实例

13. 验收标准

  • 4 张新表创建成功,DDL 幂等可重复执行
  • Flow CRUD 接口全部通过(save/publish/list/detail/archive/delete)
  • 手动启动 flow 后,start_node 正确创建 tasks 实例
  • AND 前置语义:两个前置节点都完成才触发下游
  • OR 前置语义:任一前置节点完成即触发下游
  • timeout 边:节点 deadline 到期自动触发兜底节点
  • failed 边:手动标记失败后触发兜底节点
  • count 循环终止:generation 达到 max_loops 后节点终止并触发 timeout 边
  • condition 循环终止:条件达成节点标记 completed;未达成到达 max_loops 则终止
  • 版本隔离:v1 flow 的执行在 v2 发布后不受影响
  • 并发安全:2 个回调同时到达,节点任务不重复创建
  • Quartz job 30 秒轮询有效(execution 中 pending 节点被正确评估)
  • 前端画布可拖拽节点、连线、编辑属性、发布校验
  • 执行详情树形展示正确
  • 通过 mvn clean compile 编译验证