UFO Galaxy Task Constellation 全解析用 DAG 建模与编排跨设备分布式任务流【免费下载链接】UFOUFO³: Weaving the Digital Agent Galaxy项目地址: https://gitcode.com/GitHub_Trending/uf/UFOTask Constellation任务星座是 UFO Galaxy 项目中用于捕获分布式任务执行并发与异步结构的核心抽象。它以有向无环图DAG形式对复杂工作流进行形式化建模为跨异构设备的任务提供一致的调度、容错编排与运行时动态调整能力。读完本文你将掌握 TaskStar、TaskStarLine、TaskConstellation、ConstellationEditor 四大核心组件的设计原理与完整 API 用法能够动手构建、校验、分析并动态修改自己的工作流 DAG。Task Constellation 示意图同时体现顺序依赖与并行依赖的执行拓扑核心组件四个构成星座的基本构件Task Constellation 框架由四个核心组件构成其职责分工如下源码导出见 galaxy/constellation/init.py组件职责关键能力TaskStar原子执行单元自包含的任务描述、设备指派、执行状态、依赖关系、优先级、超时与重试TaskStarLine依赖关系有向边支持无条件、仅成功、仅完成、条件四种执行语义TaskConstellationDAG 编排器完整工作流图校验、调度、动态增删改、并行度分析、序列化ConstellationEditor交互式编辑器基于命令模式的接口带 undo/redo 的安全星座操作四个组件分别对应源码文件 galaxy/constellation/task_star.py、galaxy/constellation/task_star_line.py、galaxy/constellation/task_constellation.py以及 galaxy/constellation/editor/constellation_editor.py。所有枚举定义集中在 galaxy/constellation/enums.py。形式化模型星座的数学基础一个 Task Constellation $\mathcal{C}$ 在数学上被定义为有向无环图$$ \mathcal{C} (\mathcal{T}, \mathcal{E}) $$其中 $\mathcal{T}$ 是所有TaskStar任务节点的集合$\mathcal{E}$ 是所有TaskStarLine依赖边的集合。TaskStar 的表示每个 TaskStar $t_i \in \mathcal{T}$ 封装了一份完整的任务规格$$ t_i (\text{name}_ i, \text{description}_ i, \text{target_device_id}_ i, \text{tips}_ i, \text{status}_ i, \text{dependencies}_ i) $$各字段含义name任务的简短名称description发送给设备 Agent 的自然语言任务描述target_device_id负责执行该任务的设备 Agent IDtips辅助设备 Agent 完成任务的一组提示清单status当前执行状态pending、running、completed、failed、cancelled、waiting_dependencydependencies必须先完成的前置任务 ID 集合。源码层面TaskStar 构造函数 还补充了更多可配置项task_id缺省时自动生成 UUID、device_type目标设备类型、priority默认MEDIUM、timeout秒、retry_count默认 0、task_data附加数据字典、expected_output_type与configTaskConfiguration对象用于覆盖 timeout/retry_count/priority 并合并 metadata。TaskStarLine 的表示每条 TaskStarLine $e_{i \rightarrow j} \in \mathcal{E}$ 表示从任务 $t_i$ 到任务 $t_j$ 的依赖关系。其依赖类型语义如下与 enums.py 中 DependencyType 一一对应类型行为源码判定逻辑Unconditional无条件$t_j$ 始终等待 $t_i$ 完成恒返回 TrueSuccess-only仅成功仅当 $t_i$ 成功时 $t_j$ 才继续prerequisite_result is not NoneCompletion-only仅完成$t_i$ 无论成败完成后 $t_j$ 即继续恒返回 TrueConditional条件依据用户自定义或运行时条件决定 $t_j$ 是否继续调用condition_evaluator未提供时退化为仅成功语义条件判定由 evaluate_condition() 实现值得注意的两点其一条件求值器内部抛出的异常会被捕获并返回 False不会向上传播导致编排中断其二修改dependency_type或condition_evaluator会重置满足状态is_satisfied因此在星座执行过程中修改依赖需要格外谨慎。四大关键优势1. 显式的任务排序任务依赖被显式记录在 DAG 结构中分布式执行时不存在顺序歧义保证正确性。2. 天然的并行性DAG 拓扑天然暴露可并行任务。get_ready_tasks()会返回所有满足前置依赖的待执行任务并按TaskPriorityLOW1、MEDIUM2、HIGH3、CRITICAL4降序排列见 task_constellation.py。3. 运行时动态性与静态 DAG 调度器不同Task Constellation 是可变对象任务与依赖边可以随时添加引入新的子任务或诊断任务移除剪除已完成或冗余节点修改重连依赖、更新条件、更换设备指派。这允许在不重启整个工作流的前提下实现自适应执行。4. 形式化保证DAG 表示提供了三类形式化性质无环性Acyclicity不存在循环依赖——通过 add_dependency 前的 DFS 环检测_would_create_cycle与 Kahn 拓扑排序双重保障因果一致性Causal consistency若 $t_j$ 依赖 $t_i$则 $t_i$ 必须先于 $t_j$ 完成传递依赖被完整保留并发任务之间无因果排序安全并发Safe concurrency并行执行无竞态条件——运行中任务禁止修改相关 setter 会抛出ValueError修改操作受状态检查保护。生命周期状态机星座在执行过程中会经历如下状态状态转移由 update_state() 依据全部任务状态自动推导状态说明触发条件源码逻辑CREATED星座已初始化尚未添加任务构造后无任务READY任务与依赖已配置可执行存在任务但无运行/完成/失败状态EXECUTING至少一个任务正在运行或已完成存在 RUNNING 或 COMPLETED 任务COMPLETED所有任务成功完成全部终态且无失败FAILED所有任务失败全部终态且无成功PARTIALLY_FAILED部分成功、部分失败全部终态且有失败有成功需要说明的是enums.py 中 ConstellationState 还定义了CANCELLED状态供编排层在任务被取消时使用。DAG 度量指标量化工作流并行度Task Constellation 提供多项指标用于分析工作流并行度全部实现在 task_constellation.py 中。关键路径长度Critical Path Length, $L$星座中最长的串行依赖链$$ L \max_{p \in \text{paths}} |p| $$其中 $|p|$ 为从任意根节点到任意叶子节点的路径 $p$ 的长度。源码中get_longest_path()基于拓扑序动态规划计算并返回最长路径上的任务 ID 列表get_critical_path_length_with_time()则使用真实执行时长计算以秒为单位的耗时版关键路径仅在所有任务进入终态后有效。总工作量Total Work, $W$所有任务执行时长的总和$$ W \sum_{t_i \in \mathcal{T}} \text{duration}(t_i) $$并行度比率Parallelism Ratio, $P$$$ P \frac{W}{L} $$$P 1$完全串行执行$P 1$存在并行执行空间$P$ 越高并行潜力越大。最大宽度Maximum Width同一层级可并发执行的最大任务数$$ \text{MaxWidth} \max_{\text{level}} |\text{tasks at level}| $$get_max_width()通过 BFS 层级遍历计算各层节点数取最大值。计算模式说明get_parallelism_metrics()支持两种计算模式——Node Count Mode节点计数模式执行未完成时使用任务数量每个任务计 1 单位Actual Time Mode实际时间模式所有任务进入终态后使用真实执行时长。返回值中的calculation_mode字段会明确标识当前使用哪种模式。核心操作从构建到分析DAG 构建from galaxy.constellation import TaskConstellation, TaskStar, TaskStarLine # 创建星座 constellation TaskConstellation(namemy_workflow) # 添加任务 task_a TaskStar(nametask_a, descriptionCheckout code on laptop) task_b TaskStar(nametask_b, descriptionBuild on GPU server) task_c TaskStar(nametask_c, descriptionDeploy to staging) constellation.add_task(task_a) constellation.add_task(task_b) constellation.add_task(task_c) # 添加依赖 dep_ab TaskStarLine.create_success_only( from_task_idtask_a.task_id, to_task_idtask_b.task_id, descriptionBuild depends on successful checkout ) dep_bc TaskStarLine.create_unconditional( from_task_idtask_b.task_id, to_task_idtask_c.task_id, descriptionDeploy after build ) constellation.add_dependency(dep_ab) constellation.add_dependency(dep_bc)三个便捷工厂方法create_unconditional()、create_success_only()、create_conditional()均在 task_star_line.py 中实现其中create_conditional()需要显式传入condition_evaluator回调。add_task遇到重复 task_id 会抛出ValueErroradd_dependency会校验源/目标任务存在性并通过_would_create_cycle做 DFS 环检测一旦成环同样抛出ValueError。DAG 校验# 校验结构 is_valid, errors constellation.validate_dag() if not is_valid: print(fValidation errors: {errors}) # 检查环 has_cycles constellation.has_cycle() # 获取拓扑序Kahn 算法有环时抛 ValueError order constellation.get_topological_order() print(fExecution order: {order})validate_dag()会检查是否存在环、依赖是否引用不存在的任务get_topological_order()采用 Kahn 算法结果即建议的执行顺序。并行度分析# 获取并行度指标 metrics constellation.get_parallelism_metrics() print(fCritical Path Length: {metrics[critical_path_length]}) print(fTotal Work: {metrics[total_work]}) print(fParallelism Ratio: {metrics[parallelism_ratio]}) print(fCritical Path: {metrics[critical_path_tasks]}) # 获取最大宽度 max_width constellation.get_max_width() print(fMaximum concurrent tasks: {max_width})此外get_statistics()会一次性返回星座 ID、状态、任务/依赖数量、按状态分类的任务计数、最长路径、最大宽度、并行度比率、执行时长及全部时间戳便于监控与审计。执行流编排层TaskConstellationOrchestrator见 galaxy/constellation/orchestrator/orchestrator.py通常通过以下调用推进执行constellation.start_execution() # 标记星座开始 newly_ready constellation.mark_task_completed( task_idfetch_data, successTrue, result{rows: 10000, status: success} ) # 返回因该任务完成而新就绪的任务列表mark_task_completed会遍历以该任务为源的所有依赖边逐条调用evaluate_condition()判定条件是否满足满足则从目标任务的依赖集合中移除该前置最终返回新就绪任务。TaskStar 单任务的异步执行则通过await task.execute(device_manager)调用ConstellationDeviceManager.assign_task_to_device完成超时时间缺省为 1000 秒。动态修改ConstellationEditor 安全编辑带撤销/重做的编辑from galaxy.constellation.editor import ConstellationEditor # 创建带历史记录的编辑器 editor ConstellationEditor(constellation, enable_historyTrue, max_history_size100) # 添加新诊断任务 diagnostic_task editor.create_and_add_task( task_iddiag_1, descriptionCheck server health, nameServer Health Check ) # 添加条件依赖 editor.create_and_add_dependency( from_task_idtask_b.task_id, to_task_iddiagnostic_task.task_id, dependency_typeCONDITIONAL, condition_descriptionRun diagnostic if build fails ) # 出错可撤销 if something_wrong: editor.undo() # 也可 editor.redo() 重做 # 获取可修改组件 modifiable_tasks constellation.get_modifiable_tasks() modifiable_deps constellation.get_modifiable_dependencies()ConstellationEditor 基于命令模式实现核心结构为CommandInvoker 命令注册表galaxy/constellation/editor/command_invoker.py、galaxy/constellation/editor/command_registry.py将每次操作封装为可逆命令对象从而获得完整操作历史、审计追踪、原子事务与易扩展性。它同时支持观察者模式add_observer、批量操作batch_operations、子图提取create_subgraph与星座合并merge_constellation是 Constellation Agent 构建任务工作流的主力工具。修改安全边界修改安全限制任务与依赖只有在PENDING或WAITING_DEPENDENCY状态下才可修改。运行中或已完成的任务不可修改以保证执行一致性并防止竞态条件。这一约束在源码中有多处体现TaskStar 的 name/description/tips/priority 等 setter 在RUNNING状态下抛出ValueErrorget_modifiable_tasks()只返回PENDING/WAITING_DEPENDENCY状态的任务get_modifiable_dependencies()只返回目标任务尚未开始执行的依赖边。典型工作流示例顺序工作流并行度比率1.0完全串行最大宽度1并行工作流并行度比率2.0B 与 C 可并行最大宽度2复杂工作流并行度比率约 1.67最大宽度3A 完成后 B、C、E 可并发这三种模式分别对应顺序流水线、扇出/扇入Fan-out/Fan-in与管道并行等常见组合。更完整的可运行示例线性管道、并行扇出、菱形模式含assert验证见 TaskConstellation 文档源码级测试用例可参考 tests/galaxy/constellation/。可视化多模式 DAG 展示Task Constellation 通过 galaxy/visualization/dag_visualizer.py 中的DAGVisualizer提供四种可视化模式用于监控与调试模式展示内容Overview高层结构任务数量与整体状态TopologyDAG 拓扑图任务关系与依赖Details任务明细执行时长与状态Execution实时执行流与进度追踪# 展示星座 constellation.display_dag(modeoverview) # 或 topology、details、execution交互式 Web 可视化可参见 Galaxy WebUI。组件文档与相关阅读深入各组件请查阅TaskStar — 星座中的原子执行单元含优先级分级、重试逻辑、序列化 API 参考TaskStarLine — 连接任务的依赖关系含四种依赖类型的条件求值细节TaskConstellation — 完整 DAG 编排器含调度执行、统计监控、持久化 API 参考ConstellationEditor — 交互式编辑器含命令模式、观察者、批量操作 API 参考相关文档Constellation Orchestrator — 星座如何被调度并跨设备执行Constellation Agent — Agent 如何规划与管理星座生命周期Evaluation Metrics — 监控星座性能与分析执行模式Galaxy Overview — Galaxy 整体架构与设计原则研究背景理论根基Task Constellation 模型建立在形式化 DAG 理论与分布式系统研究之上关键性质包括通过Kahn 算法进行拓扑排序获得无环性保证见 get_topological_order拓扑序保证一致的执行顺序关键路径分析用于性能优化动态图演化在不破坏一致性的前提下支持运行时变更每次变更均触发环检测与状态重算。最佳实践设计有效的星座保持任务原子性每个 TaskStar 应代表单一、明确定义的操作最小化依赖减少不必要的依赖边以最大化并行度合理选择依赖类型错误处理路径使用条件依赖如condition_evaluatorlambda result: result is None实现失败才触发的分支尽早校验执行前运行validate_dag()监控指标跟踪并行度比率以优化工作流设计。常见模式Fan-out扇出一个任务派生多个相互独立的并行任务Fan-in扇入多个并行任务汇聚到单一任务Pipeline管道顺序阶段内嵌并行任务Conditional branching条件分支用条件依赖构造错误处理路径。常见陷阱过度并行化可能压垮资源过度耦合的依赖会降低并行度忽略校验就执行修改前不检查星座状态。下一步学习 TaskStar — 原子任务执行单元探索 TaskStarLine — 依赖关系掌握 TaskConstellation — DAG 编排尝试 ConstellationEditor — 交互式编辑【免费下载链接】UFOUFO³: Weaving the Digital Agent Galaxy项目地址: https://gitcode.com/GitHub_Trending/uf/UFO创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
