add rollout=n

This commit is contained in:
鸣山 2025-07-25 14:01:05 +08:00
parent 0b95a3bd70
commit 906d3d680a
3 changed files with 173 additions and 63 deletions

View file

@ -35,6 +35,7 @@ class AppworldReactAgent:
temperature: float = 0.9,
max_interactions: int = 30,
max_response_size: int = 2000,
num_runs: int = 1,
use_experience: bool = False):
self.index: int = index
@ -44,6 +45,7 @@ class AppworldReactAgent:
self.temperature: float = temperature
self.max_interactions: int = max_interactions
self.max_response_size: int = max_response_size
self.num_runs: int = num_runs
self.use_experience: bool = use_experience
self.llm_client = OpenAI()
@ -104,39 +106,42 @@ class AppworldReactAgent:
def execute(self):
result = []
for task_index, task_id in enumerate(tqdm(self.task_ids, desc=f"ray_index={self.index}")):
with AppWorld(task_id=task_id, experiment_name=self.experiment_name) as world:
history = self.prompt_messages(world=world)
before_score = self.get_reward(world)
logger.info(f"ray_id={self.index} task_index={task_index} instruction={world.task.instruction} "
f"before_score={before_score:.4f}")
# Run each task num_runs times
for run_id in range(self.num_runs):
with AppWorld(task_id=task_id, experiment_name=f"{self.experiment_name}_run_{run_id}") as world:
history = self.prompt_messages(world=world)
before_score = self.get_reward(world)
logger.info(f"ray_id={self.index} task_index={task_index} run_id={run_id} "
f"instruction={world.task.instruction} before_score={before_score:.4f}")
for i in range(self.max_interactions):
code = self.call_llm(history)
history.append({"role": "assistant", "content": code})
for i in range(self.max_interactions):
code = self.call_llm(history)
history.append({"role": "assistant", "content": code})
output = world.execute(code)
if len(output) > self.max_response_size:
logger.warning(f"output exceed max size={len(output)}")
output = output[:self.max_response_size]
history.append({"role": "user", "content": output})
output = world.execute(code)
if len(output) > self.max_response_size:
logger.warning(f"output exceed max size={len(output)}")
output = output[:self.max_response_size]
history.append({"role": "user", "content": output})
logger.info(f"ray_id={self.index} task_index={task_index} step={i} complete~")
logger.info(f"ray_id={self.index} task_index={task_index} run_id={run_id} step={i} complete~")
if world.task_completed():
break
if world.task_completed():
break
after_score = self.get_reward(world)
uplift_score = after_score - before_score
t_result = {
"task_id": world.task_id,
"experiment_name": self.experiment_name,
"task_completed": world.task_completed(),
"before_score": before_score,
"after_score": after_score,
"uplift_score": uplift_score,
"task_history": history,
}
result.append(t_result)
after_score = self.get_reward(world)
uplift_score = after_score - before_score
t_result = {
"task_id": world.task_id,
"run_id": run_id, # Add run_id field
"experiment_name": self.experiment_name,
"task_completed": world.task_completed(),
"before_score": before_score,
"after_score": after_score,
"uplift_score": uplift_score,
"task_history": history,
}
result.append(t_result)
return result
@ -164,10 +169,10 @@ class AppworldReactAgent:
def main():
dataset_name = "train"
task_ids = load_task_ids(dataset_name)
agent = AppworldReactAgent(index=0, task_ids=task_ids[0:1], experiment_name=f"jinli_{dataset_name}")
agent = AppworldReactAgent(index=0, task_ids=task_ids[0:1], experiment_name=dataset_name, num_runs=4)
result = agent.execute()
logger.info(f"result={json.dumps(result)}")
if __name__ == "__main__":
main()
main()

View file

@ -17,7 +17,7 @@ from appworld import load_task_ids
from appworld_react_agent import AppworldReactAgent
def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, use_experience: bool=False):
def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, num_runs: int = 1, use_experience: bool = False):
experiment_name = dataset_name + "_" + experiment_suffix
path: Path = Path(f"./exp_result")
path.mkdir(parents=True, exist_ok=True)
@ -33,9 +33,12 @@ def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, use_e
if max_workers > 1:
future_list: list = []
for i in range(max_workers):
# Assign tasks to each worker, ensuring each task runs num_runs times
worker_task_ids = task_ids[i::max_workers]
actor = AppworldReactAgent.remote(index=i,
task_ids=task_ids[i::max_workers],
task_ids=worker_task_ids,
experiment_name=experiment_name,
num_runs=num_runs,
use_experience=use_experience)
future = actor.execute.remote()
future_list.append(future)
@ -50,23 +53,32 @@ def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, use_e
else:
result.append(t_result)
logger.info(f"{i + 1}/{len(task_ids)} complete")
logger.info(f"worker {i + 1}/{max_workers} complete")
dump_file()
else:
for index, task_id in enumerate(task_ids):
agent = AppworldReactAgent(index=index, task_ids=[task_id], experiment_name=experiment_name, use_experience=use_experience)
result.append(agent.execute())
agent = AppworldReactAgent(index=index,
task_ids=[task_id],
experiment_name=experiment_name,
num_runs=num_runs,
use_experience=use_experience)
task_results = agent.execute()
if isinstance(task_results, list):
result.extend(task_results)
else:
result.append(task_results)
dump_file()
def main():
max_workers = 2
max_workers = 4
num_runs = 4 # Run each task 4 times
if max_workers > 1:
ray.init(num_cpus=max_workers)
# run_agent(dataset_name="train", experiment_suffix="v2", max_workers=max_workers)
run_agent(dataset_name="dev", experiment_suffix="add-exp", max_workers=max_workers, use_experience=True)
ray.init(num_cpus=4)
# run_agent(dataset_name="train", experiment_suffix="v2", max_workers=max_workers, num_runs=num_runs)
run_agent(dataset_name="dev", experiment_suffix="add-exp", max_workers=max_workers, num_runs=num_runs, use_experience=True)
if __name__ == "__main__":
main()
main()

View file

@ -1,45 +1,138 @@
import json
from pathlib import Path
from collections import defaultdict
import pandas as pd
from loguru import logger
def calculate_best_at_k(scores: list, k: int) -> float:
"""
Calculate best@k
Divide scores into groups of size k, take the maximum value in each group,
then average these maximum values
Args:
scores: List of after_score values for all runs of a task
k: Group size
Returns:
best@k value
"""
if len(scores) % k != 0:
raise ValueError(f"Length of scores ({len(scores)}) must be divisible by k ({k})")
group_maxs = []
for i in range(0, len(scores), k):
group = scores[i:i + k]
group_maxs.append(max(group))
return sum(group_maxs) / len(group_maxs)
def get_possible_k_values(total_runs: int) -> list:
"""
Get all possible k values (factors of total_runs)
Args:
total_runs: Total number of runs
Returns:
List of k values in descending order
"""
k_values = []
for k in range(1, total_runs + 1):
if total_runs % k == 0:
k_values.append(k)
return sorted(k_values, reverse=True) # Sort from large to small
def run_exp_statistic():
path: Path = Path(f"./exp_result")
# Store results for all experiments
all_results = {}
for file in path.glob("*.jsonl"):
# Group results by task_id
task_results = defaultdict(list)
with open(file, "r") as f:
task_completed_list = []
before_score_list = []
after_score_list = []
task_success_list = []
for line in f:
if not line.strip():
continue
data = json.loads(line)
if isinstance(data, list):
for part_data in data:
task_completed_list.append(1 if part_data["task_completed"] is True else 0)
before_score_list.append(part_data["before_score"])
after_score_list.append(part_data["after_score"])
task_success_list.append(part_data["after_score"] > 0.9)
task_id = part_data["task_id"]
after_score = part_data["after_score"]
task_results[task_id].append(after_score)
else:
task_completed_list.append(1 if data["task_completed"] is True else 0)
before_score_list.append(data["before_score"])
after_score_list.append(data["after_score"])
task_success_list.append(data["after_score"] > 0.9)
task_id = data["task_id"]
after_score = data["after_score"]
task_results[task_id].append(after_score)
task_completed_ratio = sum(task_completed_list) / len(task_completed_list)
before_score_ratio = sum(before_score_list) / len(before_score_list)
after_score_ratio = sum(after_score_list) / len(after_score_list)
task_success_ratio = sum(task_success_list) / len(task_success_list)
logger.info(f"file={file} "
f"task_completed_ratio={task_completed_ratio:.2f} "
f"before_score_ratio={before_score_ratio:.2f} "
f"after_score_ratio={after_score_ratio:.2f} "
f"task_success_ratio={task_success_ratio:.2f}")
if not task_results:
logger.warning(f"No valid data found in file {file}")
continue
# Check if each task has consistent number of runs
run_counts = [len(scores) for scores in task_results.values()]
if len(set(run_counts)) > 1:
logger.warning(f"Inconsistent number of runs for different tasks in file {file}: {set(run_counts)}")
continue
num_runs = run_counts[0]
logger.info(f"File {file}: {len(task_results)} tasks, {num_runs} runs per task")
# Get all possible k values
k_values = get_possible_k_values(num_runs)
logger.info(f"Calculable best@k values: {k_values}")
# Calculate various best@k values
file_results = {"file": file.name}
for k in k_values:
best_at_k_scores = []
for task_id, scores in task_results.items():
try:
best_k_score = calculate_best_at_k(scores, k)
best_at_k_scores.append(best_k_score)
except ValueError as e:
logger.error(f"Error calculating best@{k} for task {task_id}: {e}")
continue
if best_at_k_scores:
avg_best_at_k = sum(best_at_k_scores) / len(best_at_k_scores)
file_results[f"best@{k}"] = avg_best_at_k
logger.info(f"file={file.name} best@{k}={avg_best_at_k:.4f}")
all_results[file.name] = file_results
# Create and display table
if all_results:
df = pd.DataFrame(list(all_results.values()))
df = df.set_index('file')
# Sort columns by the number in column name (best@8, best@4, best@2, best@1)
best_columns = [col for col in df.columns if col.startswith('best@')]
best_columns.sort(key=lambda x: int(x.split('@')[1]), reverse=True)
df = df[best_columns]
print("\n" + "=" * 80)
print("Experiment Results Summary Table")
print("=" * 80)
print(df.round(4))
print("=" * 80)
# Save table to CSV
output_path = path / "experiment_summary.csv"
df.to_csv(output_path)
logger.info(f"Results table saved to: {output_path}")
else:
logger.warning("No valid experiment results found")
if __name__ == "__main__":
run_exp_statistic()
run_exp_statistic()