This commit is contained in:
wty-yy
2025-12-18 13:50:23 +08:00
parent 9d509829a5
commit 1e1a04b4c0
18 changed files with 195 additions and 60 deletions

View File

@@ -8,6 +8,7 @@
@Desc : Multiprocessing Pipeline for Robogauge
'''
import yaml
import traceback
import functools
import numpy as np
from tqdm import tqdm
@@ -26,6 +27,8 @@ def run_single_process(args, data):
seed, base_mass, friction = data
local_args = deepcopy(args)
local_args.seed = seed
local_args.friction = friction
local_args.base_mass = base_mass
run_name = f"{local_args.run_name}_{seed}_baseMass{base_mass}_friction{friction}"
logger.create(
experiment_name=local_args.experiment_name,
@@ -33,8 +36,22 @@ def run_single_process(args, data):
console_output=False
)
pipeline = task_register.make_pipeline(args=local_args, create_logger=False)
log_dir = pipeline.run()
return log_dir
try:
log_dir = pipeline.run()
ret = {
'status': 'success',
'log_dir': log_dir,
'model_path': pipeline.robot_cfg.control.model_path,
}
except Exception as e:
logger.error(f"❌ Process with seed={seed}, base_mass={base_mass}, friction={friction} failed with error: {e}")
ret = {
'status': 'error',
'data': data,
'error_msg': str(e),
'traceback': traceback.format_exc()
}
return ret
class MultiPipeline:
def __init__(self, args):
@@ -43,26 +60,52 @@ class MultiPipeline:
self.frictions = args.frictions
self.base_masses = args.base_masses
self.num_processes = args.num_processes
self.model_path = None
logger.create(args.experiment_name+'_multi', args.run_name+'_multi')
def run(self):
logger.info(f"🚀 Starting Multi-Process Evaluation with {self.num_processes} processes.")
logger.info(f"🔢 Seeds: {self.seeds}, Frictions: {self.frictions}, Base masses: {self.base_masses}")
process_args = list(product(self.seeds, self.base_masses, self.frictions))
workers_data = list(product(self.seeds, self.base_masses, self.frictions))
ctx = multiprocessing.get_context('spawn')
worker_func = functools.partial(run_single_process, self.args)
result_log_dirs = []
success_flags = []
with ctx.Pool(processes=self.num_processes) as pool:
iterator = pool.imap_unordered(worker_func, process_args)
for log_dir in tqdm(iterator, total=len(process_args), desc="Evaluation"):
result_log_dirs.append(log_dir)
iterator = pool.imap_unordered(worker_func, workers_data)
for results in tqdm(iterator, total=len(workers_data), desc="Evaluation"):
success_flags.append(results['status'] == 'success')
if results['status'] == 'success':
result_log_dirs.append(results['log_dir'])
if self.model_path is None:
self.model_path = results['model_path']
else:
assert self.model_path == results['model_path'], "Model paths do not match across runs."
else:
data = results['data']
logger.error(f"❌ Process with seed={data[0]}, base_mass={data[1]}, friction={data[2]} failed with error: {results['error_msg']}")
logger.info("✅ Multi-Process Evaluation Completed.")
self.aggregate_results(result_log_dirs)
self.aggregate_results(result_log_dirs, success_flags, workers_data)
def aggregate_results(self, log_dirs):
def aggregate_results(self, log_dirs, success_flags, workers_data):
""" Process results.yaml from each log_dir """
logger.info("📊 Aggregating Results from all runs...")
summary = {'model_path': self.model_path, 'success': {}}
finish_msg = (
f"""\n{'='*20} Run Finish Summary {'='*20}\n"""
f"""{'Seed':^10}{'Base Mass':^15}{'Friction':^15}{'Status':^10}\n"""
)
for success, data in zip(success_flags, workers_data):
seed, base_mass, friction = data
status_str = "" if success else ""
finish_msg += f"{seed:^10}{base_mass:^15}{friction:^15}{status_str:^10}\n"
summary['success'][f"Seed_{seed}_BaseMass_{base_mass}_Friction_{friction}"] = True if success else False
finish_msg += f"""{'='*88}"""
logger.info(finish_msg)
all_results = []
all_yaml_paths = []
for path in log_dirs:
@@ -94,7 +137,6 @@ class MultiPipeline:
for mean_name, mean_value in means.items():
value_collections[metric][mean_name].append(float(mean_value.split(' ')[0]))
summary = {}
for metric, means in value_collections.items():
summary[metric] = {}
for mean_name, values in means.items():