From 0a3985445aebb9c45d0ff12bf0cca64c510b487a Mon Sep 17 00:00:00 2001 From: wty-yy <993660140@qq.com> Date: Sat, 27 Dec 2025 01:44:09 +0800 Subject: [PATCH] v0.1.16 --- UPDATE.md | 3 +- robogauge/tasks/pipeline/level_pipeline.py | 7 +- robogauge/tasks/pipeline/multi_pipeline.py | 13 +- robogauge/tasks/pipeline/stress_pipeline.py | 91 ++++++++----- robogauge/utils/progress_monitor.py | 137 ++++++++++++++++++++ 5 files changed, 216 insertions(+), 35 deletions(-) create mode 100644 robogauge/utils/progress_monitor.py diff --git a/UPDATE.md b/UPDATE.md index 82c8507..41c64e1 100644 --- a/UPDATE.md +++ b/UPDATE.md @@ -1,8 +1,9 @@ # UPDATE ## 20251226 ### v0.1.16 -1. 基本完成StressPipeline, 但是绘制进度信息还有问题 +1. 基本完成StressPipeline, 加入绘制进度条的线程, 其他进程通过Queue更新主进程的进度条 2. wave地形的穿模判定非常容易触发, 加入最多穿模判定重启次数为1次, 超过该次数后穿模就不再自动重启了 +3. 当全部地形等级均失败时, 需要将全部metric都按0计算 FIX Bugs: (level_results会存储到MultiPipeline下, 因为MultiPipeline不是子进程启动的, 会覆盖LevelPipeline的Logger冲突) 通过创建新的Logger实现 ## 20251225 diff --git a/robogauge/tasks/pipeline/level_pipeline.py b/robogauge/tasks/pipeline/level_pipeline.py index a2c6c59..874c6e6 100644 --- a/robogauge/tasks/pipeline/level_pipeline.py +++ b/robogauge/tasks/pipeline/level_pipeline.py @@ -11,31 +11,36 @@ import yaml from robogauge.tasks.pipeline.multi_pipeline import MultiPipeline from robogauge.utils.logger import Logger +from robogauge.utils.progress_monitor import report_progress, ProgressTypes, ProgressData level_logger = Logger() # LevelPipeline logger class LevelPipeline: - def __init__(self, args, console_output: bool = True): + def __init__(self, args, console_output=True, progress_data: ProgressData = None): self.args = args self.seeds = args.seeds self.console_output = console_output + self.progress_data = progress_data level_logger.create(args.experiment_name+'_level', args.run_name, console_output=console_output) def run(self): level_logger.info(f"🚀 Starting Level Searcher for '{self.args.experiment_name}'.") level_logger.info(f"🔢 Seeds: {self.seeds}") + report_progress(self.progress_data, ProgressTypes.INIT, total=10, desc="🔍 Searching Max Level") # binary search levels l, r = 0, 10 all_level_results = {} while l < r: level = (l + r + 1) // 2 + report_progress(self.progress_data, ProgressTypes.DESC, desc=f"🔍 Testing Level {level}") all_success, results = self.test_level(level) all_level_results[level] = results if all_success: l = level else: r = level - 1 + report_progress(self.progress_data, ProgressTypes.UPDATE, value=1) level = l level_results = all_level_results.get(l, { 'model_path': results['model_path'], diff --git a/robogauge/tasks/pipeline/multi_pipeline.py b/robogauge/tasks/pipeline/multi_pipeline.py index 49868bc..93cad20 100644 --- a/robogauge/tasks/pipeline/multi_pipeline.py +++ b/robogauge/tasks/pipeline/multi_pipeline.py @@ -23,6 +23,7 @@ from robogauge.tasks.pipeline.base_pipeline import BasePipeline from robogauge.utils.task_register import task_register from robogauge.utils.logger import Logger from robogauge.utils.process_utils import NoDaemonPool +from robogauge.utils.progress_monitor import report_progress, ProgressTypes, ProgressData multi_logger = Logger() # MultiPipeline logger @@ -63,11 +64,13 @@ def run_single_process(args, data): return ret class MultiPipeline: - def __init__(self, args, console_output: bool = True): + def __init__(self, args, console_output=True, progress_data: ProgressData = None): self.args = args self.seeds = args.seeds self.frictions = args.frictions self.base_masses = args.base_masses + self.console_output = console_output + self.progress_data = progress_data self.num_processes = args.num_processes self.static_info = {} multi_logger.create(args.experiment_name+'_multi', args.run_name+'_multi', console_output=console_output) @@ -83,16 +86,22 @@ class MultiPipeline: multi_logger.info(f"🔢 Seeds: {self.seeds}, Frictions: {self.frictions}, Base masses: {self.base_masses}") workers_data = list(product(self.seeds, self.base_masses, self.frictions)) + report_progress(self.progress_data, ProgressTypes.INIT, total=len(workers_data), desc="🚀 MultiPipeline Executing") + ctx = multiprocessing.get_context('spawn') worker_func = functools.partial(run_single_process, self.args) results_list = [] with NoDaemonPool(processes=self.num_processes, context=ctx) as pool: iterator = pool.imap_unordered(worker_func, workers_data) - for results in tqdm(iterator, total=len(workers_data), desc="Evaluation"): + bar = iterator + if self.console_output: + bar = tqdm(iterator, total=len(workers_data), desc="Evaluation") + for results in bar: results_list.append(results) self.add_static_info('model_path', results['model_path']) self.add_static_info('terrain_name', results['results']['terrain_name']) self.add_static_info('terrain_level', results['results']['terrain_level']) + report_progress(self.progress_data, ProgressTypes.UPDATE, value=1) if results['status'] != 'success': data = results['data'] multi_logger.error(f"❌ Process with seed={data[0]}, base_mass={data[1]}, friction={data[2]} failed with error: {results['error_msg']}") diff --git a/robogauge/tasks/pipeline/stress_pipeline.py b/robogauge/tasks/pipeline/stress_pipeline.py index fefcd53..a52ee0c 100644 --- a/robogauge/tasks/pipeline/stress_pipeline.py +++ b/robogauge/tasks/pipeline/stress_pipeline.py @@ -8,6 +8,7 @@ @Desc : Stress Pipeline for Robogauge ''' import yaml +import traceback import functools import numpy as np from tqdm import tqdm @@ -18,45 +19,57 @@ from collections import defaultdict from robogauge.utils.logger import Logger from robogauge.utils.process_utils import NoDaemonPool +from robogauge.utils.progress_monitor import report_progress, ProgressTypes, start_progress_monitor_thread, ProgressData from robogauge.tasks.pipeline import MultiPipeline, LevelPipeline from robogauge.tasks.gauge.gauge_configs.terrain_levels_config import TerrainSearchLevelsConfig stress_logger = Logger() # StressPipeline logger -def run_pipeline(args, data): +def run_pipeline(args, progress_queue, data): args = deepcopy(args) + task_id = data['task_id'] search = data['search_max_level'] + task_label = f"[{data['terrain_name']}]" + progress_data = ProgressData( + task_id=task_id, + msg_prefix=task_label + ' ', + progress_queue=progress_queue + ) if search is True: + task_label += f" M:{data['base_mass']} F:{data['friction']}" + progress_data.msg_prefix = task_label + ' ' args.friction = data['friction'] args.frictions = [data['friction']] args.base_mass = data['base_mass'] args.base_masses = [data['base_mass']] args.task_name = f"{data['task_robot_model']}.{data['terrain_name']}" - args.experiment_name = f"{args.experiment_name}_{data['terrain_name']}_baseMass{data['base_mass']}_friction{data['friction']}" + args.experiment_name = f"{args.experiment_name}_{data['terrain_name']}_M{data['base_mass']}_F{data['friction']}" + + level, level_results = LevelPipeline(args, console_output=False, progress_data=progress_data).run() + if level == 0: # no valid level found + report_progress(progress_data, ProgressTypes.FINISH, desc=f"❌ Failed (Lv 0)") + results = { + 'success': False, + 'results': level_results, + 'data': data, + 'level': 0, + } + return results + + report_progress(progress_data, ProgressTypes.RESET, total=0, desc=f"✅ Found Lv {level} -> Running") else: + level = None # flat terrain args.task_name = f"{data['task_robot_model']}.{data['terrain_name']}" args.experiment_name = f"{args.experiment_name}_{data['terrain_name']}" - level = None # flat terrain - if search: - level, results = LevelPipeline(args, console_output=False).run() - - if level == 0: # no valid level found - results = { - 'success': False, - 'results': results, - 'data': data, - 'level': 0, - } - else: - args.level = level - results = { - 'success': True, - 'results': MultiPipeline(args, console_output=False).run(), - 'data': data, - 'level': level, - } - + args.level = level + results = { + 'success': True, + 'results': MultiPipeline(args, console_output=False, progress_data=progress_data).run(), + 'data': data, + 'level': level, + } + report_progress(progress_data, ProgressTypes.FINISH, desc=f"✅ Done (Lv {level})") return results class StressPipeline: @@ -81,9 +94,6 @@ class StressPipeline: terrain_names = self.args.stress_terrain_names stress_logger.info(f"🌄 Stress Test Terrain Names: {terrain_names}") - ctx = multiprocessing.get_context('spawn') - worker_func = functools.partial(run_pipeline, self.args) - ### Build worker data ### workers_data = [] terrain_search_levels_config = TerrainSearchLevelsConfig() @@ -109,13 +119,26 @@ class StressPipeline: else: workers_data.append(data) + ### Start progress monitor ### + progress_queue, monitor_thread = start_progress_monitor_thread(len(workers_data)) + for i, data in enumerate(workers_data): + data['task_id'] = i + ### Run and collect results ### + ctx = multiprocessing.get_context('spawn') + worker_func = functools.partial(run_pipeline, self.args, progress_queue) results_list = [] - with NoDaemonPool(processes=self.num_processes, context=ctx) as pool: - iterator = pool.imap_unordered(worker_func, workers_data) - for results in tqdm(iterator, total=len(workers_data), desc="Stress Benchmark"): - results_list.append(results) - self.add_static_info('model_path', results['results'].pop('model_path', None)) + try: + with NoDaemonPool(processes=self.num_processes, context=ctx) as pool: + iterator = pool.imap_unordered(worker_func, workers_data) + for results in iterator: + results_list.append(results) + self.add_static_info('model_path', results['results'].pop('model_path', None)) + except Exception as e: + stress_logger.error(f"❌ Stress benchmark encountered an error: {e}, {traceback.format_exc()}") + finally: + progress_queue.put(None) # Stop the progress monitor thread + monitor_thread.join() stress_logger.info("✅ Stress Benchmark Completed.") stress_results = self.aggregate_results(results_list) @@ -144,6 +167,7 @@ class StressPipeline: summary = {**self.static_info, 'summary': {}} value_collections = defaultdict(lambda: defaultdict(list)) + zero_terrain_count = 0 for result in all_results: terrain_name = result['data']['terrain_name'] terrain_level = result['level'] # None, 0, 1, ..., 10 @@ -151,7 +175,11 @@ class StressPipeline: if terrain_level is not None: key += f'_{terrain_level}' key += f'_baseMass{result["data"]["base_mass"]}_friction{result["data"]["friction"]}' - summary[key] = result['results'] if terrain_level != 0 else None + if terrain_level == 0: + summary[key] = None + zero_terrain_count += 1 + continue + summary[key] = result['results'] for metric, means in result['results']['summary'].items(): for mean_name, mean_value in means.items(): @@ -160,6 +188,7 @@ class StressPipeline: for metric, means in value_collections.items(): summary['summary'][metric] = {} for mean_name, values in means.items(): + values.extend([0.0] * zero_terrain_count) # include zero terrains summary['summary'][metric][mean_name] = f"{float(np.mean(values)):.4f} ± {float(np.std(values)):.4f}" save_path = stress_logger.log_dir / "stress_benchmark_results.yaml" diff --git a/robogauge/utils/progress_monitor.py b/robogauge/utils/progress_monitor.py new file mode 100644 index 0000000..45201d6 --- /dev/null +++ b/robogauge/utils/progress_monitor.py @@ -0,0 +1,137 @@ +# -*- coding: utf-8 -*- +''' +@File : progress_monitor.py +@Time : 2025/12/27 00:10:07 +@Author : wty-yy +@Version : 1.0 +@Blog : https://wty-yy.github.io/ +@Desc : Centralized progress monitoring for multiprocessing tasks using tqdm +''' +from tqdm import tqdm +import multiprocessing +from typing import Dict +from threading import Thread +from dataclasses import dataclass + +class ProgressTypes: + INIT = 'init' # Init progress (set total, desc) + UPDATE = 'update' # Update progress value (set value) + DESC = 'desc' # Update description text only (set desc) + RESET = 'reset' # Reset progress bar (set total, desc) + FINISH = 'finish' # Mark completion (set desc) + ERROR = 'error' # Mark error (set desc) + +@dataclass +class ProgressData: + progress_queue: multiprocessing.Queue + task_id: int = 0 + msg_prefix: str = '' + +def report_progress(progress_data: ProgressData, msg_type, value=None, desc=None, total=None): + """ + Assistant function to report progress to the main process. + Args: + queue: multiprocessing.Queue + task_id: int, line number of the task + msg_type: ProgressTypes + value: int, value to update (for UPDATE type) + desc: str, description text (for DESC, INIT, FINISH types) + total: int, total value (for INIT, RESET types) + """ + if progress_data is None: + return + queue = progress_data.progress_queue + task_id = progress_data.task_id + msg_prefix = progress_data.msg_prefix + if desc is not None: + desc = msg_prefix + desc + queue.put((task_id, msg_type, {'value': value, 'desc': desc, 'total': total})) + +class ProgressMonitor: + def __init__(self, total_rows): + self.total_rows = total_rows + self.bars: Dict[int, tqdm] = {} + + def listener_loop(self, queue): + """ + Run in a separate thread in the main process to consume the Queue and update tqdm. + """ + # print(f"Monitor started for {self.total_rows} tasks...") + for i in range(self.total_rows): + self.bars[i] = tqdm( + total=100, + position=i, + desc=f"Task {i} Pending...", + bar_format="{l_bar}{bar}| {n_fmt}/{total_fmt} [{elapsed}]", + leave=True + ) + + active_tasks = self.total_rows + + while active_tasks > 0: + record = queue.get() + if record is None: # Poison pill signal + break + + task_id, msg_type, data = record + if task_id not in self.bars: + continue + + bar = self.bars[task_id] + + if msg_type == ProgressTypes.INIT: + total = data['total'] + desc = data['desc'] + bar.reset(total=total) + bar.set_description(desc) + bar.refresh() + + elif msg_type == ProgressTypes.UPDATE: + val = data['value'] + bar.update(val) + + elif msg_type == ProgressTypes.DESC: + desc = data['desc'] + if desc: + bar.set_description(desc) + + elif msg_type == ProgressTypes.RESET: + # Scenario: Level search finished, starting MultiPipeline, reset progress bar + total = data['total'] + desc = data['desc'] + bar.reset(total=total) + bar.set_description(desc) + bar.refresh() + + elif msg_type == ProgressTypes.FINISH: + desc = data['desc'] + if desc: + bar.set_description(desc) + bar.refresh() + # Note: We do not close the bar here to keep it displayed until all tasks are finished and closed together. + active_tasks -= 1 + + elif msg_type == ProgressTypes.ERROR: + desc = data['desc'] + bar.set_description(desc) + bar.refresh() + active_tasks -= 1 + + # After all tasks are finished, close all bars + for bar in self.bars.values(): + bar.close() + +def start_progress_monitor_thread(total_rows): + """ + Create and start a ProgressMonitor thread. Queue is returned for reporting progress. + Args: + total_rows: int, number of tasks to monitor + Returns: + queue: multiprocessing.Queue + monitor_thread: threading.Thread + """ + queue = multiprocessing.Manager().Queue() + monitor = ProgressMonitor(total_rows) + monitor_thread = Thread(target=monitor.listener_loop, args=(queue,)) + monitor_thread.start() + return queue, monitor_thread