优先级: P1 预计工时: 后端 5d + 前端 3d 状态: 设计已确认(待写实现计划) 日期: 2026-09-19
现有任务系统(tasks 表)已支持单前置(prerequisite_task_id)、循环(repeat_type)、父链(parent_task_id),但缺少:
这些能力需要新增上层编排模块,不修改现有 tasks 表结构,复用现有 TaskService 的创建/进度/积分体系。
| 决策项 | 结论 |
|---|---|
| 驱动场景 | 家庭/成长任务的条件编排 |
| 最大复杂度 | 含环 DAG |
| 前置语义 | AND(全部) + OR(任一)均支持 |
| 未完成语义 | 超时未按时完成 + 主动标记失败,均支持 |
| 流定义方式 | UI 可视化编排(jsPlumb 画布) |
| 节点与 tasks 关系 | 节点执行时创建 tasks 实例,复用现有进度/积分体系 |
| 循环终止 | 计数终止 + 条件终止两者都支持 |
| 调度模型 | 引入 Quartz 独立轮询,不依赖 @Scheduled |
| 内部回调范围 | 仅编排节点任务回调,普通任务零开销 |
| 画布选型 | jsPlumb(Element UI 兼容,vue2 有现成封装) |
| 版本管理 | 执行实例快照旧版本,新执行用新版本 |
| 并发保护 | 乐观锁 + 幂等检查 |
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 |
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 |
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 |
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) — 防止重复创建。
每个节点的完整配置(存储在 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
}
对每个 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 启动时直接触发。
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 出边| 触发来源 | 节点状态变化 | 边类型 | 处理 |
|---|---|---|---|
| deadline 到期且未完成 | pending → timeout |
edge_type=timeout |
创建目标节点任务实例 |
| 用户手动标记放弃 | pending → failed |
edge_type=failed |
创建目标节点任务实例 |
| 异常(如任务模板不存在) | pending → failed |
edge_type=failed |
同上 |
节点标记为 skipped 时(flow 终止/暂停),不触发任何出边。
当节点触发时:
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
三个方法仅对属于编排节点的任务生效(先查 task_orchestration_node_instances.task_id = taskId,找不到则直接返回,零开销):
| 方法 | 时机 | 效果 |
|---|---|---|
onTaskCompleted(Long taskId) |
任务完成/审核通过 | 节点 → completed,评估下游 |
onTaskFailed(Long taskId) |
用户主动放弃/异常 | 节点 → failed,触发 failed 边 |
onTaskTimeout(Long taskId) |
deadline 到期未完成 | 节点 → timeout,触发 timeout 边 |
evaluateFlow(executionId) 对同一个 execution_id 加分布式锁(Redis 或数据库行锁),避免 Quartz 轮询与回调同时触发导致重复创建节点任务(execution_id, node_id, generation) 作为最终防线@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)
}
}
除 Quartz 轮询外,以下节点状态变更即时触发评估:
onTaskCompleted / onTaskFailed 回调(事件驱动,无需等待下一个轮询周期)node/fail 接口调用统一 @PostMapping,路由前缀 /api/orchestration。
| 接口 | 说明 | 请求体关键字段 |
|---|---|---|
/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 |
| 接口 | 说明 | 请求体关键字段 |
|---|---|---|
/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 |
| 接口 | 说明 | 请求体关键字段 |
|---|---|---|
/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 |
见 §5.5。
/orchestration — 编排流管理首页(列表 + 新建入口)/orchestration/edit/:id — 编排画布编辑器/orchestration/execution/:id — 执行详情只读页admin_task_templates 读取)is_start_node = true| 场景 | 处理策略 |
|---|---|
| 节点引用的任务模板被删除 | 发布校验拦截;运行时节点标记 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 循环) |
OrchestrationEngine.evaluateFlow():覆盖 AND/OR 语义、无前置节点、多入边混合在 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 建表语句。
| 不做 | 原因 |
|---|---|
| 家庭挑战/五维打卡合并进编排 | 已有设计文档明确不合并 |
| 前端小程序端编排编辑 | 编排编辑只放 cfc-web 管理端;小程序端只做执行结果展示(复用 tasks 列表) |
| 节点级精确 cron 定时 | 只用 flow 级 schedule_cron,节点触发完全由依赖边决定 |
| 跨家庭共享执行实例 | 执行实例严格绑定 family,系统级 flow 仅用于多个家庭复用定义 |
| 流市场/模板市场 | 后续迭代 |
| 节点级并行执行(同时创建多个任务) | 当前版本不支持,一个节点同一时刻只有一个任务实例 |
mvn clean compile 编译验证