This commit is contained in:
wty-yy
2025-12-27 01:44:09 +08:00
parent dad65a4e27
commit 0a3985445a
5 changed files with 216 additions and 35 deletions

View File

@@ -1,8 +1,9 @@
# UPDATE # UPDATE
## 20251226 ## 20251226
### v0.1.16 ### v0.1.16
1. 基本完成StressPipeline, 但是绘制进度信息还有问题 1. 基本完成StressPipeline, 加入绘制进度条的线程, 其他进程通过Queue更新主进程的进度条
2. wave地形的穿模判定非常容易触发, 加入最多穿模判定重启次数为1次, 超过该次数后穿模就不再自动重启了 2. wave地形的穿模判定非常容易触发, 加入最多穿模判定重启次数为1次, 超过该次数后穿模就不再自动重启了
3. 当全部地形等级均失败时, 需要将全部metric都按0计算
FIX Bugs: (level_results会存储到MultiPipeline下, 因为MultiPipeline不是子进程启动的, 会覆盖LevelPipeline的Logger冲突) 通过创建新的Logger实现 FIX Bugs: (level_results会存储到MultiPipeline下, 因为MultiPipeline不是子进程启动的, 会覆盖LevelPipeline的Logger冲突) 通过创建新的Logger实现
## 20251225 ## 20251225

View File

@@ -11,31 +11,36 @@ import yaml
from robogauge.tasks.pipeline.multi_pipeline import MultiPipeline from robogauge.tasks.pipeline.multi_pipeline import MultiPipeline
from robogauge.utils.logger import Logger from robogauge.utils.logger import Logger
from robogauge.utils.progress_monitor import report_progress, ProgressTypes, ProgressData
level_logger = Logger() # LevelPipeline logger level_logger = Logger() # LevelPipeline logger
class LevelPipeline: 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.args = args
self.seeds = args.seeds self.seeds = args.seeds
self.console_output = console_output self.console_output = console_output
self.progress_data = progress_data
level_logger.create(args.experiment_name+'_level', args.run_name, console_output=console_output) level_logger.create(args.experiment_name+'_level', args.run_name, console_output=console_output)
def run(self): def run(self):
level_logger.info(f"🚀 Starting Level Searcher for '{self.args.experiment_name}'.") level_logger.info(f"🚀 Starting Level Searcher for '{self.args.experiment_name}'.")
level_logger.info(f"🔢 Seeds: {self.seeds}") level_logger.info(f"🔢 Seeds: {self.seeds}")
report_progress(self.progress_data, ProgressTypes.INIT, total=10, desc="🔍 Searching Max Level")
# binary search levels # binary search levels
l, r = 0, 10 l, r = 0, 10
all_level_results = {} all_level_results = {}
while l < r: while l < r:
level = (l + r + 1) // 2 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_success, results = self.test_level(level)
all_level_results[level] = results all_level_results[level] = results
if all_success: if all_success:
l = level l = level
else: else:
r = level - 1 r = level - 1
report_progress(self.progress_data, ProgressTypes.UPDATE, value=1)
level = l level = l
level_results = all_level_results.get(l, { level_results = all_level_results.get(l, {
'model_path': results['model_path'], 'model_path': results['model_path'],

View File

@@ -23,6 +23,7 @@ from robogauge.tasks.pipeline.base_pipeline import BasePipeline
from robogauge.utils.task_register import task_register from robogauge.utils.task_register import task_register
from robogauge.utils.logger import Logger from robogauge.utils.logger import Logger
from robogauge.utils.process_utils import NoDaemonPool from robogauge.utils.process_utils import NoDaemonPool
from robogauge.utils.progress_monitor import report_progress, ProgressTypes, ProgressData
multi_logger = Logger() # MultiPipeline logger multi_logger = Logger() # MultiPipeline logger
@@ -63,11 +64,13 @@ def run_single_process(args, data):
return ret return ret
class MultiPipeline: 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.args = args
self.seeds = args.seeds self.seeds = args.seeds
self.frictions = args.frictions self.frictions = args.frictions
self.base_masses = args.base_masses self.base_masses = args.base_masses
self.console_output = console_output
self.progress_data = progress_data
self.num_processes = args.num_processes self.num_processes = args.num_processes
self.static_info = {} self.static_info = {}
multi_logger.create(args.experiment_name+'_multi', args.run_name+'_multi', console_output=console_output) 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}") 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)) 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') ctx = multiprocessing.get_context('spawn')
worker_func = functools.partial(run_single_process, self.args) worker_func = functools.partial(run_single_process, self.args)
results_list = [] results_list = []
with NoDaemonPool(processes=self.num_processes, context=ctx) as pool: with NoDaemonPool(processes=self.num_processes, context=ctx) as pool:
iterator = pool.imap_unordered(worker_func, workers_data) 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) results_list.append(results)
self.add_static_info('model_path', results['model_path']) self.add_static_info('model_path', results['model_path'])
self.add_static_info('terrain_name', results['results']['terrain_name']) self.add_static_info('terrain_name', results['results']['terrain_name'])
self.add_static_info('terrain_level', results['results']['terrain_level']) self.add_static_info('terrain_level', results['results']['terrain_level'])
report_progress(self.progress_data, ProgressTypes.UPDATE, value=1)
if results['status'] != 'success': if results['status'] != 'success':
data = results['data'] 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']}") multi_logger.error(f"❌ Process with seed={data[0]}, base_mass={data[1]}, friction={data[2]} failed with error: {results['error_msg']}")

View File

@@ -8,6 +8,7 @@
@Desc : Stress Pipeline for Robogauge @Desc : Stress Pipeline for Robogauge
''' '''
import yaml import yaml
import traceback
import functools import functools
import numpy as np import numpy as np
from tqdm import tqdm from tqdm import tqdm
@@ -18,45 +19,57 @@ from collections import defaultdict
from robogauge.utils.logger import Logger from robogauge.utils.logger import Logger
from robogauge.utils.process_utils import NoDaemonPool 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.pipeline import MultiPipeline, LevelPipeline
from robogauge.tasks.gauge.gauge_configs.terrain_levels_config import TerrainSearchLevelsConfig from robogauge.tasks.gauge.gauge_configs.terrain_levels_config import TerrainSearchLevelsConfig
stress_logger = Logger() # StressPipeline logger stress_logger = Logger() # StressPipeline logger
def run_pipeline(args, data): def run_pipeline(args, progress_queue, data):
args = deepcopy(args) args = deepcopy(args)
task_id = data['task_id']
search = data['search_max_level'] 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: 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.friction = data['friction']
args.frictions = [data['friction']] args.frictions = [data['friction']]
args.base_mass = data['base_mass'] args.base_mass = data['base_mass']
args.base_masses = [data['base_mass']] args.base_masses = [data['base_mass']]
args.task_name = f"{data['task_robot_model']}.{data['terrain_name']}" 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: else:
level = None # flat terrain
args.task_name = f"{data['task_robot_model']}.{data['terrain_name']}" args.task_name = f"{data['task_robot_model']}.{data['terrain_name']}"
args.experiment_name = f"{args.experiment_name}_{data['terrain_name']}" args.experiment_name = f"{args.experiment_name}_{data['terrain_name']}"
level = None # flat terrain args.level = level
if search: results = {
level, results = LevelPipeline(args, console_output=False).run() 'success': True,
'results': MultiPipeline(args, console_output=False, progress_data=progress_data).run(),
if level == 0: # no valid level found 'data': data,
results = { 'level': level,
'success': False, }
'results': results, report_progress(progress_data, ProgressTypes.FINISH, desc=f"✅ Done (Lv {level})")
'data': data,
'level': 0,
}
else:
args.level = level
results = {
'success': True,
'results': MultiPipeline(args, console_output=False).run(),
'data': data,
'level': level,
}
return results return results
class StressPipeline: class StressPipeline:
@@ -81,9 +94,6 @@ class StressPipeline:
terrain_names = self.args.stress_terrain_names terrain_names = self.args.stress_terrain_names
stress_logger.info(f"🌄 Stress Test Terrain Names: {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 ### ### Build worker data ###
workers_data = [] workers_data = []
terrain_search_levels_config = TerrainSearchLevelsConfig() terrain_search_levels_config = TerrainSearchLevelsConfig()
@@ -109,13 +119,26 @@ class StressPipeline:
else: else:
workers_data.append(data) 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 ### ### Run and collect results ###
ctx = multiprocessing.get_context('spawn')
worker_func = functools.partial(run_pipeline, self.args, progress_queue)
results_list = [] results_list = []
with NoDaemonPool(processes=self.num_processes, context=ctx) as pool: try:
iterator = pool.imap_unordered(worker_func, workers_data) with NoDaemonPool(processes=self.num_processes, context=ctx) as pool:
for results in tqdm(iterator, total=len(workers_data), desc="Stress Benchmark"): iterator = pool.imap_unordered(worker_func, workers_data)
results_list.append(results) for results in iterator:
self.add_static_info('model_path', results['results'].pop('model_path', None)) 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_logger.info("✅ Stress Benchmark Completed.")
stress_results = self.aggregate_results(results_list) stress_results = self.aggregate_results(results_list)
@@ -144,6 +167,7 @@ class StressPipeline:
summary = {**self.static_info, 'summary': {}} summary = {**self.static_info, 'summary': {}}
value_collections = defaultdict(lambda: defaultdict(list)) value_collections = defaultdict(lambda: defaultdict(list))
zero_terrain_count = 0
for result in all_results: for result in all_results:
terrain_name = result['data']['terrain_name'] terrain_name = result['data']['terrain_name']
terrain_level = result['level'] # None, 0, 1, ..., 10 terrain_level = result['level'] # None, 0, 1, ..., 10
@@ -151,7 +175,11 @@ class StressPipeline:
if terrain_level is not None: if terrain_level is not None:
key += f'_{terrain_level}' key += f'_{terrain_level}'
key += f'_baseMass{result["data"]["base_mass"]}_friction{result["data"]["friction"]}' 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 metric, means in result['results']['summary'].items():
for mean_name, mean_value in means.items(): for mean_name, mean_value in means.items():
@@ -160,6 +188,7 @@ class StressPipeline:
for metric, means in value_collections.items(): for metric, means in value_collections.items():
summary['summary'][metric] = {} summary['summary'][metric] = {}
for mean_name, values in means.items(): 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}" 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" save_path = stress_logger.log_dir / "stress_benchmark_results.yaml"

View File

@@ -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