Files
NexQuant/rdagent/app/data_science/loop.py
T

358 lines
16 KiB
Python
Raw Normal View History

import shutil
import subprocess
from datetime import datetime
from pathlib import Path
from typing import Any, Optional, Union
import fire
from rdagent.app.data_science.conf import DS_RD_SETTING
from rdagent.components.coder.data_science.ensemble import EnsembleCoSTEER
2025-02-13 22:20:17 +08:00
from rdagent.components.coder.data_science.ensemble.exp import EnsembleTask
from rdagent.components.coder.data_science.feature import FeatureCoSTEER
2025-02-13 22:20:17 +08:00
from rdagent.components.coder.data_science.feature.exp import FeatureTask
from rdagent.components.coder.data_science.model import ModelCoSTEER
2025-02-13 22:20:17 +08:00
from rdagent.components.coder.data_science.model.exp import ModelTask
2025-04-03 13:00:11 +08:00
from rdagent.components.coder.data_science.pipeline import PipelineCoSTEER
from rdagent.components.coder.data_science.pipeline.exp import PipelineTask
from rdagent.components.coder.data_science.raw_data_loader import DataLoaderCoSTEER
2025-02-13 22:20:17 +08:00
from rdagent.components.coder.data_science.raw_data_loader.exp import DataLoaderTask
from rdagent.components.coder.data_science.share.doc import DocDev
from rdagent.components.coder.data_science.workflow import WorkflowCoSTEER
2025-02-13 22:20:17 +08:00
from rdagent.components.coder.data_science.workflow.exp import WorkflowTask
from rdagent.components.workflow.conf import BasePropSetting
from rdagent.components.workflow.rd_loop import RDLoop
from rdagent.core.conf import RD_AGENT_SETTINGS
from rdagent.core.exception import CoderError, RunnerError
from rdagent.core.proposal import ExperimentFeedback
from rdagent.core.scenario import Scenario
from rdagent.core.utils import import_class
from rdagent.log import rdagent_logger as logger
from rdagent.scenarios.data_science.dev.feedback import DSExperiment2Feedback
from rdagent.scenarios.data_science.dev.runner import DSCoSTEERRunner
from rdagent.scenarios.data_science.experiment.experiment import DSExperiment
from rdagent.scenarios.data_science.proposal.exp_gen import DSExpGen, DSTrace
2025-05-09 15:38:25 +08:00
from rdagent.scenarios.data_science.proposal.exp_gen.ckp_select import (
BackJumpCKPSelector,
LatestCKPSelector,
SOTAJumpCKPSelector,
)
from rdagent.scenarios.data_science.proposal.exp_gen.idea_pool import DSKnowledgeBase
2025-05-09 15:38:25 +08:00
from rdagent.scenarios.data_science.proposal.exp_gen.sota_exp_select import (
AutoSOTAexpSelector,
BestValidSelector,
GlobalSOTASelector,
)
from rdagent.scenarios.kaggle.kaggle_crawler import download_data
2025-05-09 15:38:25 +08:00
CKP_SELECTOR_NAME_MAP = {
"latest": LatestCKPSelector,
"sota_jump": SOTAJumpCKPSelector,
"back_jump": BackJumpCKPSelector,
}
SOTA_EXP_SELECTOR_NAME_MAP = {
"global_sota": GlobalSOTASelector,
"auto_sota": AutoSOTAexpSelector,
"best_valid_sota": BestValidSelector,
}
class DataScienceRDLoop(RDLoop):
skip_loop_error = (CoderError, RunnerError)
def __init__(self, PROP_SETTING: BasePropSetting):
logger.log_object(PROP_SETTING.competition, tag="competition")
scen: Scenario = import_class(PROP_SETTING.scen)(PROP_SETTING.competition)
# 1) task generation from scratch
# self.scratch_gen: tuple[HypothesisGen, Hypothesis2Experiment] = DummyHypothesisGen(scen),
# 2) task generation from a complete solution
# self.exp_gen: ExpGen = import_class(PROP_SETTING.exp_gen)(scen)
2025-05-09 15:38:25 +08:00
# self.ckp_selector = CKP_SELECTOR_NAME_MAP[DS_RD_SETTING.selector_name]()
# self.sota_exp_selector = SOTA_EXP_SELECTOR_NAME_MAP[DS_RD_SETTING.sota_exp_selector_name]()
self.ckp_selector = import_class(PROP_SETTING.selector_name)()
self.sota_exp_selector = import_class(PROP_SETTING.sota_exp_selector_name)()
self.exp_gen = import_class(PROP_SETTING.hypothesis_gen)(scen)
# coders
self.data_loader_coder = DataLoaderCoSTEER(scen)
self.feature_coder = FeatureCoSTEER(scen)
self.model_coder = ModelCoSTEER(scen)
self.ensemble_coder = EnsembleCoSTEER(scen)
self.workflow_coder = WorkflowCoSTEER(scen)
2025-04-03 13:00:11 +08:00
self.pipeline_coder = PipelineCoSTEER(scen)
self.runner = DSCoSTEERRunner(scen)
if DS_RD_SETTING.enable_doc_dev:
self.docdev = DocDev(scen)
# self.summarizer: Experiment2Feedback = import_class(PROP_SETTING.summarizer)(scen)
# logger.log_object(self.summarizer, tag="summarizer")
if DS_RD_SETTING.enable_knowledge_base and DS_RD_SETTING.knowledge_base_version == "v1":
knowledge_base = DSKnowledgeBase(
path=DS_RD_SETTING.knowledge_base_path, idea_pool_json_path=DS_RD_SETTING.idea_pool_json_path
)
self.trace = DSTrace(scen=scen, knowledge_base=knowledge_base)
else:
self.trace = DSTrace(scen=scen)
self.summarizer = DSExperiment2Feedback(scen)
super(RDLoop, self).__init__()
def direct_exp_gen(self, prev_out: dict[str, Any]):
2025-05-09 15:38:25 +08:00
# set the SOTA experiment to submit
sota_exp_to_submit = self.sota_exp_selector.get_sota_exp_to_submit(self.trace)
self.trace.set_sota_exp_to_submit(sota_exp_to_submit)
# set the checkpoint to start from
2025-04-09 09:42:30 +08:00
selection = self.ckp_selector.get_selection(self.trace)
exp = self.exp_gen.gen(self.trace, selection)
2025-01-23 16:12:22 +08:00
logger.log_object(exp)
# FIXME: this is for LLM debug webapp, remove this when the debugging is done.
logger.log_object(exp, tag="debug_exp_gen")
return exp
def coding(self, prev_out: dict[str, Any]):
exp = prev_out["direct_exp_gen"]
for tasks in exp.pending_tasks_list:
exp.sub_tasks = tasks
2025-03-20 19:15:37 +08:00
with logger.tag(f"{exp.sub_tasks[0].__class__.__name__}"):
if isinstance(exp.sub_tasks[0], DataLoaderTask):
exp = self.data_loader_coder.develop(exp)
elif isinstance(exp.sub_tasks[0], FeatureTask):
exp = self.feature_coder.develop(exp)
elif isinstance(exp.sub_tasks[0], ModelTask):
exp = self.model_coder.develop(exp)
elif isinstance(exp.sub_tasks[0], EnsembleTask):
exp = self.ensemble_coder.develop(exp)
elif isinstance(exp.sub_tasks[0], WorkflowTask):
exp = self.workflow_coder.develop(exp)
2025-04-03 13:00:11 +08:00
elif isinstance(exp.sub_tasks[0], PipelineTask):
exp = self.pipeline_coder.develop(exp)
2025-03-20 19:15:37 +08:00
else:
raise NotImplementedError(f"Unsupported component in DataScienceRDLoop: {exp.hypothesis.component}")
exp.sub_tasks = []
2025-01-23 16:12:22 +08:00
logger.log_object(exp)
return exp
def running(self, prev_out: dict[str, Any]):
exp: DSExperiment = prev_out["coding"]
if exp.is_ready_to_run():
2025-01-20 12:19:21 +08:00
new_exp = self.runner.develop(exp)
2025-01-23 16:12:22 +08:00
logger.log_object(new_exp)
exp = new_exp
if DS_RD_SETTING.enable_doc_dev:
self.docdev.develop(exp)
return exp
def feedback(self, prev_out: dict[str, Any]) -> ExperimentFeedback:
"""
Assumption:
- If we come to feedback phase, the previous development steps are successful.
"""
exp: DSExperiment = prev_out["running"]
2025-04-03 13:00:11 +08:00
if self.trace.next_incomplete_component() is None or DS_RD_SETTING.coder_on_whole_pipeline:
# we have alreadly completed components in previous trace. So current loop is focusing on a new proposed idea.
# So we need feedback for the proposal.
feedback = self.summarizer.generate_feedback(exp, self.trace)
else:
# Otherwise, it is on drafting stage, don't need complicated feedbacks.
feedback = ExperimentFeedback(
reason=f"{exp.hypothesis.component} is completed.",
decision=True,
)
2025-01-23 16:12:22 +08:00
logger.log_object(feedback)
return feedback
def record(self, prev_out: dict[str, Any]):
2025-04-09 09:42:30 +08:00
# set the DAG parent for the trace
self.trace.sync_dag_parent_and_hist()
e = prev_out.get(self.EXCEPTION_KEY, None)
if e is None:
self.trace.hist.append((prev_out["running"], prev_out["feedback"]))
else:
self.trace.hist.append(
(
prev_out["direct_exp_gen"] if isinstance(e, CoderError) else prev_out["coding"],
ExperimentFeedback.from_exception(e),
)
)
if self.trace.sota_experiment() is None:
if DS_RD_SETTING.coder_on_whole_pipeline:
# check if feedback is not generated
if len(self.trace.hist) >= DS_RD_SETTING.coding_fail_reanalyze_threshold:
recent_hist = self.trace.hist[-DS_RD_SETTING.coding_fail_reanalyze_threshold :]
if all(isinstance(fb.exception, (CoderError, RunnerError)) for _, fb in recent_hist):
new_scen = self.trace.scen
if hasattr(new_scen, "reanalyze_competition_description"):
logger.info(
"Reanalyzing the competition description after three consecutive coding failures."
)
new_scen.reanalyze_competition_description()
self.trace.scen = new_scen
else:
logger.info("Can not reanalyze the competition description.")
elif len(self.trace.hist) >= DS_RD_SETTING.consecutive_errors:
# if {in inital/drafting stage} and {tried enough times}
for _, fb in self.trace.hist[-DS_RD_SETTING.consecutive_errors :]:
if fb:
break # any success will stop restarting.
else: # otherwise restart it
logger.error("Consecutive errors reached the limit. Dumping trace.")
logger.log_object(self.trace, tag="trace before restart")
self.trace = DSTrace(scen=self.trace.scen, knowledge_base=self.trace.knowledge_base)
logger.log_object(self.trace, tag="trace")
logger.log_object(self.trace.sota_experiment(), tag="SOTA experiment")
2025-04-29 09:30:45 +08:00
if DS_RD_SETTING.enable_knowledge_base and DS_RD_SETTING.knowledge_base_version == "v1":
logger.log_object(self.trace.knowledge_base, tag="knowledge_base")
self.trace.knowledge_base.dump()
if (
DS_RD_SETTING.enable_log_archive
and DS_RD_SETTING.log_archive_path is not None
and Path(DS_RD_SETTING.log_archive_path).is_dir()
):
start_archive_datetime = datetime.now()
logger.info(f"Archiving log and workspace folder after loop {self.loop_idx}")
mid_log_tar_path = (
Path(
DS_RD_SETTING.log_archive_temp_path
if DS_RD_SETTING.log_archive_temp_path
else DS_RD_SETTING.log_archive_path
)
/ "mid_log.tar"
)
mid_workspace_tar_path = (
Path(
DS_RD_SETTING.log_archive_temp_path
if DS_RD_SETTING.log_archive_temp_path
else DS_RD_SETTING.log_archive_path
)
/ "mid_workspace.tar"
)
subprocess.run(["tar", "-cf", str(mid_log_tar_path), "-C", (Path().cwd() / "log"), "."], check=True)
# remove all files and folders in the workspace except for .py, .md, and .csv files to avoid large workspace dump
for workspace_id in Path(RD_AGENT_SETTINGS.workspace_path).iterdir():
for file_and_folder in workspace_id.iterdir():
if file_and_folder.is_dir():
shutil.rmtree(file_and_folder)
elif file_and_folder.is_file() and file_and_folder.suffix not in [".py", ".md", ".csv"]:
file_and_folder.unlink()
subprocess.run(
["tar", "-cf", str(mid_workspace_tar_path), "-C", (RD_AGENT_SETTINGS.workspace_path), "."], check=True
)
if DS_RD_SETTING.log_archive_temp_path is not None:
shutil.move(mid_log_tar_path, Path(DS_RD_SETTING.log_archive_path) / "mid_log.tar")
mid_log_tar_path = Path(DS_RD_SETTING.log_archive_path) / "mid_log.tar"
shutil.move(mid_workspace_tar_path, Path(DS_RD_SETTING.log_archive_path) / "mid_workspace.tar")
mid_workspace_tar_path = Path(DS_RD_SETTING.log_archive_path) / "mid_workspace.tar"
shutil.copy(
mid_log_tar_path, Path(DS_RD_SETTING.log_archive_path) / "mid_log_bak.tar"
) # backup when upper code line is killed when running
shutil.copy(
mid_workspace_tar_path, Path(DS_RD_SETTING.log_archive_path) / "mid_workspace_bak.tar"
) # backup when upper code line is killed when running
self.timer.add_duration(datetime.now() - start_archive_datetime)
@classmethod
def load(
cls,
path: Union[str, Path],
output_path: Optional[Union[str, Path]] = None,
do_truncate: bool = False,
replace_timer: bool = True,
) -> "LoopBase":
session = super().load(path, output_path, do_truncate, replace_timer)
2025-04-29 09:30:45 +08:00
logger.log_object(DS_RD_SETTING.competition, tag="competition") # NOTE: necessary to make mle_summary work.
if DS_RD_SETTING.enable_knowledge_base and DS_RD_SETTING.knowledge_base_version == "v1":
session.trace.knowledge_base = DSKnowledgeBase(
path=DS_RD_SETTING.knowledge_base_path, idea_pool_json_path=DS_RD_SETTING.idea_pool_json_path
)
return session
def dump(self, path: str | Path) -> None:
"""
Since knowledge_base is big and we don't want to dump it every time
So we remove it from the trace before dumping and restore it after.
"""
backup_knowledge_base = None
if self.trace.knowledge_base is not None:
backup_knowledge_base = self.trace.knowledge_base
self.trace.knowledge_base = None
super().dump(path)
if backup_knowledge_base is not None:
self.trace.knowledge_base = backup_knowledge_base
def main(
path=None,
output_path=None,
step_n=None,
loop_n=None,
competition="bms-molecular-translation",
do_truncate=True,
timeout=None,
replace_timer=True,
2025-04-29 09:30:45 +08:00
exp_gen_cls: str | None = None,
):
"""
Parameters
----------
path :
path like `$LOG_PATH/__session__/1/0_propose`. It indicates that we restore the state that after finish the step 0 in loop 1
output_path :
path like `$LOG_PATH`. It indicates that where we want to save our session and log information.
step_n :
How many steps to run; if None, it will run forever until error or KeyboardInterrupt
loop_n :
How many loops to run; if None, it will run forever until error or KeyboardInterrupt
- if current loop is incomplete, it will be counted as the first loop for completion.
- if both step_n and loop_n are provided, the process will stop as soon as either condition is met.
competition :
do_truncate :
If set to True, the logger will truncate the future log messages by calling `logger.storage.truncate`.
2025-04-29 09:30:45 +08:00
replace_timer :
If session is loaded, should we replace the timer with session.timer
exp_gen_cls :
When we have different stages, we can replace the exp_gen with the new proposal
Auto R&D Evolving loop for models in a Kaggle scenario.
You can continue running session by
.. code-block:: bash
dotenv run -- python rdagent/app/data_science/loop.py [--competition titanic] $LOG_PATH/__session__/1/0_propose --step_n 1 # `step_n` is a optional parameter
rdagent kaggle --competition playground-series-s4e8 # You are encouraged to use this one.
"""
if competition is not None:
DS_RD_SETTING.competition = competition
2025-05-06 16:00:13 +08:00
if not DS_RD_SETTING.competition:
logger.error("Please specify competition name.")
2025-05-06 16:00:13 +08:00
if path is None:
kaggle_loop = DataScienceRDLoop(DS_RD_SETTING)
else:
kaggle_loop = DataScienceRDLoop.load(path, output_path, do_truncate, replace_timer)
2025-04-29 09:30:45 +08:00
# replace exp_gen if we have new class
if exp_gen_cls is not None:
kaggle_loop.exp_gen = import_class(exp_gen_cls)(kaggle_loop.exp_gen.scen)
kaggle_loop.run(step_n=step_n, loop_n=loop_n, all_duration=timeout)
if __name__ == "__main__":
fire.Fire(main)