# 任务编排系统设计 **优先级:** 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`): ```json { "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 ```java @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()` 新增迁移: ```sql 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` 编译验证