UFO Galaxy 会话指标观测器(SessionMetricsObserver)深度解析:星座执行的性能度量与统计分析
【免费下载链接】UFOUFO³: Weaving the Digital Agent Galaxy项目地址: https://gitcode.com/GitHub_Trending/uf/UFO
SessionMetricsObserver 是 UFO Galaxy 框架观察者体系(Observer System)中的核心指标组件,负责在星座(Constellation)执行过程中收集任务执行时间、星座生命周期、结构修改等三类原始度量,并在查询时计算成功率、平均耗时、修改分布等统计摘要。本文以 metrics_observer.md 为主线,结合 base_observer.py、events.py、galaxy_session.py 等源码实现,系统讲解该观察者的架构、事件处理链路、统计计算逻辑、集成方式与实战用法,帮助读者理解并复用在 Galaxy 星座执行中做性能评估与瓶颈分析的关键能力。
一、定位与设计目标:为什么需要 Metrics Observer
UFO Galaxy 采用事件驱动的观察者模式来解耦星座执行过程中的事件生产与消费。在 overview.md 描述的观察者体系中,Orchestrator、Agent等事件发布者只需向全局EventBus发布事件,而观测者(Observer)则按需订阅并处理这些事件。SessionMetricsObserver正是专门负责性能度量和统计的观察者。
其核心设计目标可归纳为四点:
- 性能跟踪(Performance Tracking)——度量任务与星座的执行耗时;
- 成功率监控(Success Rate Monitoring)——追踪完成与失败的数量和比率;
- 修改分析(Modification Analytics)——监控星座结构的动态变更;
- 统计摘要(Statistical Summaries)——在查询时聚合计算各类统计指标,供研究、评测与优化使用。
从源码结构看,SessionMetricsObserver实现了IEventObserver抽象接口(定义于 events.py),与ConstellationProgressObserver一同位于 base_observer.py 中,并在galaxy/session/observers/__init__.py中对外导出。它是评估 Galaxy 性能、定位瓶颈、分析星座修改模式的关键基础设施。
二、架构总览:事件源、事件总线与输出
原文档以 mermaid 图描述了指标的流转路径:Orchestrator发布任务事件、Agent发布星座事件,二者汇入全局EventBus;EventBus通知SessionMetricsObserver,观察者内部通过任务事件处理器与星座事件处理器将原始数据写入Metrics Storage;查询时由Statistics Computer计算统计摘要,最终输出到result.json与日志。
对照源码,该图的每个环节都有真实对应:
- 事件总线:
EventBus是 events.py 中的单例,维护_observers(按事件类型索引)与_all_observers两组订阅集合。发布事件时通过asyncio.gather并发通知所有相关观察者,单个观察者的异常不会阻断其他观察者。 - 事件分派:
SessionMetricsObserver.on_event依据事件类型分派——TaskEvent走任务处理分支,ConstellationEvent走星座处理分支(见 base_observer.py)。 - 统计计算:
get_metrics()在返回时即时调用_compute_task_statistics()、_compute_constellation_statistics()、_compute_modification_statistics()三个私有方法完成聚合(见 base_observer.py)。
观察者在体系中的位置
在观察者体系中,SessionMetricsObserver与其余四个观察者并行工作、互不干扰(详见 overview.md):
| 观察者 | 源码文件 | 核心职责 |
|---|---|---|
| SessionMetricsObserver | base_observer.py | 性能度量:耗时、成功率、修改统计 |
| ConstellationProgressObserver | base_observer.py | 任务进度跟踪、向 Agent 投递完成事件 |
| DAGVisualizationObserver | dag_visualization_observer.py | 星座拓扑实时可视化 |
| ConstellationModificationSynchronizer | constellation_sync_observer.py | 修改同步与竞态防护 |
| AgentOutputObserver | agent_output_observer.py | Agent 响应与动作的实时展示 |
三、采集的指标体系
SessionMetricsObserver将指标划分为任务指标、星座指标、修改指标三大类。前两类在事件到达时实时累加(Real-time),第三类及衍生指标在get_metrics()调用时计算(Computed)。
3.1 任务指标(Task Metrics)
跟踪单个任务的执行过程:
| Metric | 描述 | 计算时机 |
|---|---|---|
| task_count | 启动的任务总数 | Real-time |
| completed_tasks | 成功完成的任务数 | Real-time |
| failed_tasks | 失败的任务数 | Real-time |
| total_execution_time | 所有任务执行耗时之和 | Real-time |
| task_timings | 字典:task_id → {start, end, duration} | Real-time |
| success_rate | completed / total 任务数 | Computed |
| failure_rate | failed / total 任务数 | Computed |
| average_task_duration | 每个任务的平均执行时间 | Computed |
| min_task_duration | 最快任务执行时间 | Computed |
| max_task_duration | 最慢任务执行时间 | Computed |
3.2 星座指标(Constellation Metrics)
监控星座的完整生命周期:
| Metric | 描述 | 计算时机 |
|---|---|---|
| constellation_count | 处理的星座总数 | Real-time |
| completed_constellations | 成功完成的星座数 | Real-time |
| failed_constellations | 失败的星座数 | Real-time |
| total_constellation_time | 星座总执行时间 | Real-time |
| constellation_timings | 字典:constellation_id → 时序数据 | Real-time |
| constellation_success_rate | completed / total 星座数 | Computed |
| average_constellation_duration | 平均星座执行时间 | Computed |
| min_constellation_duration | 最快星座 | Computed |
| max_constellation_duration | 最慢星座 | Computed |
| average_tasks_per_constellation | 每个星座的平均任务数 | Computed |
3.3 修改指标(Modification Metrics)
跟踪星座的结构性变更(动态 DAG 调整):
| Metric | 描述 | 计算时机 |
|---|---|---|
| constellation_modifications | 字典:constellation_id → 修改记录列表 | Real-time |
| total_modifications | 修改总次数 | Computed |
| constellations_modified | 发生过修改的星座数 | Computed |
| average_modifications_per_constellation | 平均每个星座的修改次数 | Computed |
| max_modifications_for_single_constellation | 单星座最大修改次数 | Computed |
| most_modified_constellation | 被修改最多的星座 ID | Computed |
| modification_types_breakdown | 按修改类型统计的次数分布 | Computed |
3.4 内部度量存储结构
上述原始指标由构造函数初始化的字典统一维护(见 base_observer.py),与文档描述完全一致:
self.metrics: Dict[str, Any] = { "session_id": session_id, # Task metrics "task_count": 0, "completed_tasks": 0, "failed_tasks": 0, "total_execution_time": 0.0, "task_timings": {}, # task_id -> {start, end, duration} # Constellation metrics "constellation_count": 0, "completed_constellations": 0, "failed_constellations": 0, "total_constellation_time": 0.0, "constellation_timings": {}, # constellation_id -> timing data # Modification tracking "constellation_modifications": {} # constellation_id -> [modifications] }四、初始化与订阅:接入事件总线
4.1 构造与订阅
SessionMetricsObserver只需两个参数即可创建:必填的session_id(会话唯一标识)和可选的logger(不传则自动创建默认 Logger)。
from galaxy.session.observers import SessionMetricsObserver import logging # Create metrics observer metrics_observer = SessionMetricsObserver( session_id="galaxy_session_20231113", logger=logging.getLogger(__name__) ) # Subscribe to event bus from galaxy.core.events import get_event_bus event_bus = get_event_bus() event_bus.subscribe(metrics_observer)构造参数:
| 参数 | 类型 | 必填 | 描述 |
|---|---|---|---|
session_id | str | 是 | 会话唯一标识,用于区分不同执行批次 |
logger | logging.Logger | 否 | Logger 实例,为 None 时自动创建默认实例 |
需要注意:EventBus.subscribe(observer)不传event_types时订阅全部事件(见 events.py)。由于on_event内部通过isinstance判断事件类型,SessionMetricsObserver会过滤掉AgentEvent、DeviceEvent等无关事件,仅处理TaskEvent与ConstellationEvent。若希望更精确,也可用subscribe(observer, {EventType.TASK_COMPLETED, ...})形式订阅指定事件类型以降低通知开销。
4.2 在 GalaxySession 中的自动装配
在实际项目中,观察者通常由GalaxySession自动装配,而非手动创建。在 galaxy_session.py 的_setup_observers()方法中:
# Metrics observer for performance tracking self._metrics_observer = SessionMetricsObserver( session_id=f"galaxy_session_{self._id}", logger=self.logger ) self._observers.append(self._metrics_observer)会话执行结束后,get_metrics()的完整结果被写入self._session_results["metrics"](见 galaxy_session.py),可通过session.session_results属性直接访问(见 galaxy_session.py)。这是将指标与整次会话结果一体化交付的标准路径。
五、事件处理链路:从事件到指标
SessionMetricsObserver.on_event是唯一的入口,根据事件类型分派到六个内部处理方法:
async def on_event(self, event: Event) -> None: if isinstance(event, TaskEvent): await self._handle_task_event(event) elif isinstance(event, ConstellationEvent): await self._handle_constellation_event(event)5.1 任务事件处理
任务事件覆盖TASK_STARTED、TASK_COMPLETED、TASK_FAILED三种类型(EventType枚举定义见 events.py):
对应的处理逻辑(见 base_observer.py):
def _handle_task_started(self, event: TaskEvent) -> None: """Handle TASK_STARTED event.""" self.metrics["task_count"] += 1 self.metrics["task_timings"][event.task_id] = {"start": event.timestamp} def _handle_task_completed(self, event: TaskEvent) -> None: """Handle TASK_COMPLETED event.""" self.metrics["completed_tasks"] += 1 if event.task_id in self.metrics["task_timings"]: duration = ( event.timestamp - self.metrics["task_timings"][event.task_id]["start"] ) self.metrics["task_timings"][event.task_id]["duration"] = duration self.metrics["task_timings"][event.task_id]["end"] = event.timestamp self.metrics["total_execution_time"] += duration def _handle_task_failed(self, event: TaskEvent) -> None: """Handle TASK_FAILED event.""" self.metrics["failed_tasks"] += 1 # Also calculate duration for failed tasks if event.task_id in self.metrics["task_timings"]: duration = ( event.timestamp - self.metrics["task_timings"][event.task_id]["start"] ) self.metrics["task_timings"][event.task_id]["duration"] = duration self.metrics["total_execution_time"] += duration要点:失败任务同样计入耗时统计——即使任务失败,其从开始到失败的时间差也会写入task_timings并累加进total_execution_time,保证耗时统计不丢失失败样本。TaskEvent的task_id、timestamp、status、result、error字段定义于 events.py。
5.2 星座事件处理
星座生命周期事件包括CONSTELLATION_STARTED、CONSTELLATION_COMPLETED、CONSTELLATION_MODIFIED(CONSTELLATION_FAILED枚举存在,观察者当前未为其设置独立分支,失败星座的统计依赖生命周期事件与状态字段)。
def _handle_constellation_started(self, event: ConstellationEvent) -> None: """Handle CONSTELLATION_STARTED event.""" self.metrics["constellation_count"] += 1 constellation_id = event.constellation_id constellation = event.data.get("constellation") # Store initial statistics self.metrics["constellation_timings"][constellation_id] = { "start_time": event.timestamp, "initial_statistics": ( constellation.get_statistics() if constellation else {} ), "processing_start_time": event.data.get("processing_start_time"), "processing_end_time": event.data.get("processing_end_time"), "processing_duration": event.data.get("processing_duration"), } def _handle_constellation_completed(self, event: ConstellationEvent) -> None: """Handle CONSTELLATION_COMPLETED event.""" self.metrics["completed_constellations"] += 1 constellation_id = event.constellation_id constellation = event.data.get("constellation") duration = ( event.timestamp - self.metrics["constellation_timings"][constellation_id]["start_time"] if constellation_id in self.metrics["constellation_timings"] else None ) if constellation_id in self.metrics["constellation_timings"]: self.metrics["constellation_timings"][constellation_id].update({ "end_time": event.timestamp, "duration": duration, "final_statistics": ( constellation.get_statistics() if constellation else {} ), })值得注意的实现细节:
- 星座开始事件会同时记录
initial_statistics(初始快照),完成事件则记录final_statistics(终态快照),两者对比即可还原星座从创建到完成的演化; - 事件
data中若携带processing_start_time/processing_end_time/processing_duration,会原样透传存储,便于与 Agent 的处理阶段耗时对齐; constellation.get_statistics()来自 task_constellation.py,返回total_tasks、total_dependencies、task_status_counts、longest_path_length、max_width以及并行度指标(L、W、P)等结构化统计,是后续计算"每个星座平均任务数"的数据来源。
5.3 修改跟踪(Modification Tracking)
星座的动态修改是 Galaxy 的进阶能力:Agent 在任务完成后可增删任务、调整依赖,从而重构图谱。_handle_constellation_modified借助VisualizationChangeDetector做细粒度变更检测:
def _handle_constellation_modified(self, event: ConstellationEvent) -> None: """Handle CONSTELLATION_MODIFIED event.""" constellation_id = event.constellation_id if constellation_id not in self.metrics["constellation_modifications"]: self.metrics["constellation_modifications"][constellation_id] = [] if hasattr(event, "data") and event.data: old_constellation = event.data.get("old_constellation") new_constellation = event.data.get("new_constellation") changes = None if old_constellation and new_constellation: changes = VisualizationChangeDetector.calculate_constellation_changes( old_constellation, new_constellation ) modification_record = { "timestamp": event.timestamp, "modification_type": event.data.get("modification_type", "unknown"), "on_task_id": event.data.get("on_task_id", []), "changes": changes, "new_statistics": ( new_constellation.get_statistics() if new_constellation else {} ), "processing_start_time": event.data.get("processing_start_time"), "processing_end_time": event.data.get("processing_end_time"), "processing_duration": event.data.get("processing_duration"), } self.metrics["constellation_modifications"][constellation_id].append( modification_record )VisualizationChangeDetector.calculate_constellation_changes定义于 change_detector.py,返回结构化的变更明细:
added_tasks/removed_tasks/modified_tasks:通过新旧星座的 task_id 集合差集与属性比对得出;added_dependencies/removed_dependencies/modified_dependencies:依赖以from_task_id->to_task_id字符串形式描述,同样先做集合差集、再做属性比对;modification_type:根据变更内容归纳为constellation_created、任务/依赖增删改等类型。
每条修改记录还携带触发修改的目标任务(on_task_id)、修改时间戳与修改后的星座统计快照,为"哪次修改发生在哪个任务完成后、带来了什么结构变化"提供了完整可追溯的审计链。
六、统计计算:get_metrics() 与三类统计摘要
get_metrics()是观察者的对外查询接口,返回原始度量与三类计算统计的组合字典:
def get_metrics(self) -> Dict[str, Any]: """Get collected metrics with computed statistics.""" metrics = self.metrics.copy() metrics["task_statistics"] = self._compute_task_statistics() metrics["constellation_statistics"] = self._compute_constellation_statistics() metrics["modification_statistics"] = self._compute_modification_statistics() return metrics返回值包含:全部原始指标(计数、时序等)+task_statistics+constellation_statistics+modification_statistics。
6.1 任务统计(Task Statistics)
{ "total_tasks": 10, "completed_tasks": 8, "failed_tasks": 2, "success_rate": 0.8, "failure_rate": 0.2, "average_task_duration": 2.5, "min_task_duration": 0.5, "max_task_duration": 5.2, "total_task_execution_time": 25.0 }源码实现要点(见 base_observer.py):durations从task_timings中提取所有带duration字段的记录;success_rate、failure_rate在task_count > 0时计算,否则返回0.0,规避除零异常。
6.2 星座统计(Constellation Statistics)
{ "total_constellations": 1, "completed_constellations": 1, "failed_constellations": 0, "success_rate": 1.0, "average_constellation_duration": 30.5, "min_constellation_duration": 30.5, "max_constellation_duration": 30.5, "total_constellation_time": 30.5, "average_tasks_per_constellation": 10.0 }实现要点(见 base_observer.py):average_tasks_per_constellation并非简单的任务总数除以星座数,而是遍历每个星座的initial_statistics["total_tasks"]求和后除以含该字段的星座数——以初始任务快照为分子口径,避免后续动态修改造成统计失真。
6.3 修改统计(Modification Statistics)
{ "total_modifications": 3, "constellations_modified": 1, "average_modifications_per_constellation": 3.0, "max_modifications_for_single_constellation": 3, "most_modified_constellation": "const_123", "modifications_per_constellation": { "const_123": 3 }, "modification_types_breakdown": { "add_tasks": 2, "modify_dependencies": 1 } }实现要点(见 base_observer.py):modification_types_breakdown对每条记录的modification_type字段做计数聚合(缺失时归入"unknown");most_modified_constellation通过max(..., key=lambda x: x[1])找出修改次数最多的星座 ID,无修改时返回None。
6.4 典型读取示例
# After constellation execution metrics = metrics_observer.get_metrics() # Access task statistics print(f"Total tasks: {metrics['task_statistics']['total_tasks']}") print(f"Success rate: {metrics['task_statistics']['success_rate']:.2%}") print(f"Avg duration: {metrics['task_statistics']['average_task_duration']:.2f}s") # Access constellation statistics print(f"Total constellations: {metrics['constellation_statistics']['total_constellations']}") print(f"Avg tasks per constellation: {metrics['constellation_statistics']['average_tasks_per_constellation']:.1f}") # Access modification statistics print(f"Total modifications: {metrics['modification_statistics']['total_modifications']}") print(f"Modification types: {metrics['modification_statistics']['modification_types_breakdown']}")七、完整使用示例
示例 1:基础指标采集
在星座执行前后订阅观察者并读取汇总结果:
import asyncio from galaxy.core.events import get_event_bus from galaxy.session.observers import SessionMetricsObserver async def collect_metrics(): """Collect and display metrics for constellation execution.""" # Create and subscribe metrics observer metrics_observer = SessionMetricsObserver(session_id="demo_session") event_bus = get_event_bus() event_bus.subscribe(metrics_observer) # Execute constellation (orchestrator will publish events) await orchestrator.execute_constellation(constellation) # Retrieve metrics metrics = metrics_observer.get_metrics() # Display summary print("\n=== Execution Summary ===") print(f"Session: {metrics['session_id']}") print(f"Tasks: {metrics['task_count']} total, " f"{metrics['completed_tasks']} completed, " f"{metrics['failed_tasks']} failed") print(f"Total execution time: {metrics['total_execution_time']:.2f}s") # Display task statistics task_stats = metrics['task_statistics'] print(f"\nTask Success Rate: {task_stats['success_rate']:.1%}") print(f"Average Task Duration: {task_stats['average_task_duration']:.2f}s") print(f"Fastest Task: {task_stats['min_task_duration']:.2f}s") print(f"Slowest Task: {task_stats['max_task_duration']:.2f}s") # Clean up event_bus.unsubscribe(metrics_observer) asyncio.run(collect_metrics())示例 2:性能分析(定位瓶颈)
利用task_timings排序找出最慢任务,并结合修改统计分析 Agent 的动态调整行为:
def analyze_performance(metrics_observer: SessionMetricsObserver): """Analyze performance metrics and identify bottlenecks.""" metrics = metrics_observer.get_metrics() task_timings = metrics['task_timings'] # Find slowest tasks sorted_tasks = sorted( task_timings.items(), key=lambda x: x[1].get('duration', 0), reverse=True ) print("\n=== Top 5 Slowest Tasks ===") for task_id, timing in sorted_tasks[:5]: duration = timing.get('duration', 0) print(f"{task_id}: {duration:.2f}s") # Analyze modification patterns mod_stats = metrics['modification_statistics'] if mod_stats['total_modifications'] > 0: print(f"\n=== Modification Analysis ===") print(f"Total Modifications: {mod_stats['total_modifications']}") print(f"Average per Constellation: " f"{mod_stats['average_modifications_per_constellation']:.1f}") print(f"Most Modified: {mod_stats['most_modified_constellation']}") print("\nModification Types:") for mod_type, count in mod_stats['modification_types_breakdown'].items(): print(f" {mod_type}: {count}")示例 3:导出 JSON 供离线分析
将统计结果序列化为 JSON 文件,便于后续研究或报表生成:
import json from pathlib import Path def export_metrics(metrics_observer: SessionMetricsObserver, output_path: str): """Export metrics to JSON file for analysis.""" metrics = metrics_observer.get_metrics() # Convert to JSON-serializable format output_data = { "session_id": metrics["session_id"], "task_statistics": metrics["task_statistics"], "constellation_statistics": metrics["constellation_statistics"], "modification_statistics": metrics["modification_statistics"], "raw_metrics": { "task_count": metrics["task_count"], "completed_tasks": metrics["completed_tasks"], "failed_tasks": metrics["failed_tasks"], "total_execution_time": metrics["total_execution_time"], "constellation_count": metrics["constellation_count"], } } # Write to file output_file = Path(output_path) output_file.parent.mkdir(parents=True, exist_ok=True) with open(output_file, 'w') as f: json.dump(output_data, f, indent=2) print(f"Metrics exported to: {output_file}")八、最佳实践
1. Session ID 命名
使用描述性会话 ID 便于后续跨批次对比分析:
# ✅ Good: Descriptive session ID session_id = f"galaxy_session_{task_type}_{timestamp}" # ❌ Bad: Generic session ID session_id = "session_1"2. 无论成败都导出指标
将指标导出放进finally,即使执行异常也能保留诊断数据:
try: await orchestrator.execute_constellation(constellation) finally: # Always export metrics, even if execution failed metrics = metrics_observer.get_metrics() export_metrics(metrics, "results/metrics.json")3. 长会话的内存管理
task_timings与constellation_timings会随任务量线性增长,长跑会话处理完指标后应清理:
# After processing metrics metrics_observer.metrics["task_timings"].clear() metrics_observer.metrics["constellation_timings"].clear()4. 生命周期管理
与其余观察者一致,会话结束或执行异常后应统一unsubscribe,避免观察者持续接收事件造成内存泄漏与多余开销(模式详见 overview.md 中的 Observer Lifecycle Management 章节)。
5. 复用现有装配而非重复创建
在GalaxySession场景下,观察者已被 galaxy_session.py 自动创建并订阅,应通过session_results["metrics"]读取,而非在外部重复实例化。
九、测试与验证
仓库中为观察者体系提供了多层次的测试覆盖,可作为理解行为与验证自定义扩展的参考:
- test_session_observers.py —— 会话观察者的功能测试;
- test_modular_observers.py 与 test_observer_modular_structure.py —— 观察者模块化结构的单元测试;
- test_event_system.py ——
EventBus、EventType与订阅/发布机制的测试; - test_galaxy_framework_summary.py —— 框架级汇总测试,覆盖会话结果与指标输出。
这些测试一方面验证了事件分派、计数累加与统计计算的正确性,另一方面也展示了如何在自定义场景中组合EventBus与观察者进行可复现的指标采集。
十、总结
SessionMetricsObserver是 UFO Galaxy 观察者体系中负责性能度量和统计分析的关键组件,通过订阅全局EventBus,在星座执行期间被动采集任务与星座的生命周期事件,实时维护计数与时序数据,并在查询时聚合出三类统计摘要:
- 收集完整的性能指标:任务数、完成/失败数、执行耗时;
- 跟踪任务与星座的执行时间及星座初始/终态统计快照;
- 监控星座结构修改模式,借助
VisualizationChangeDetector记录细粒度的 DAG 变更明细; - 计算成功率、平均/最值耗时、每星座平均任务数、修改类型分布等统计摘要;
- 导出结构化数据(JSON / session_results)供性能评估、瓶颈定位与科研分析。
配合事件系统核心文档(event_system.md)与观察者体系总览(overview.md),即可完整理解并复用这套指标采集能力。在涉及星座执行与编排的深入场景中,可进一步参阅 constellation_orchestrator 目录 下的编排实现与文档。
【免费下载链接】UFOUFO³: Weaving the Digital Agent Galaxy项目地址: https://gitcode.com/GitHub_Trending/uf/UFO
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考