mirror of
https://github.com/NicolasBohn/NexQuant.git
synced 2026-08-01 17:37:43 +00:00
fix: fix the problems weights bug (#898)
* fix the problems weights bug * refactor: remove DSExpGen * update problems weights calculation * update problems weights calculation * remove the selection parameter from exp_gen * v2 support draft * v3 also support decomposition * make the identify_problems an independent function * fix minor bug * reformat * rename exp_num to weighted_exp_num * add the set_current_selection before the exp_gen when merging * reformat * fix wrong selection * refactor: drop selection arg from ExpGen.gen and DS merge generators --------- Co-authored-by: Xu <v-xuminrui@microsoft.com> Co-authored-by: Young <afe.young@gmail.com> Co-authored-by: Xu Yang <xuyang1@microsoft.com>
This commit is contained in:
@@ -14,7 +14,7 @@ class DataScienceBasePropSetting(KaggleBasePropSetting):
|
||||
scen: str = "rdagent.scenarios.data_science.scen.KaggleScen"
|
||||
"""Scenario class for data mining model"""
|
||||
|
||||
hypothesis_gen: str = "rdagent.scenarios.data_science.proposal.exp_gen.DSExpGen"
|
||||
hypothesis_gen: str = "rdagent.scenarios.data_science.proposal.exp_gen.proposal.DSProposalV2ExpGen"
|
||||
"""Hypothesis generation class"""
|
||||
|
||||
## Workflow Related
|
||||
|
||||
@@ -1,298 +1,9 @@
|
||||
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
|
||||
from rdagent.components.coder.data_science.ensemble.exp import EnsembleTask
|
||||
from rdagent.components.coder.data_science.feature import FeatureCoSTEER
|
||||
from rdagent.components.coder.data_science.feature.exp import FeatureTask
|
||||
from rdagent.components.coder.data_science.model import ModelCoSTEER
|
||||
from rdagent.components.coder.data_science.model.exp import ModelTask
|
||||
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
|
||||
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
|
||||
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
|
||||
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
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.sota_exp_select import (
|
||||
AutoSOTAexpSelector,
|
||||
BestValidSelector,
|
||||
GlobalSOTASelector,
|
||||
)
|
||||
from rdagent.scenarios.kaggle.kaggle_crawler import download_data
|
||||
|
||||
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)
|
||||
|
||||
# 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)
|
||||
|
||||
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]):
|
||||
|
||||
# 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
|
||||
selection = self.ckp_selector.get_selection(self.trace)
|
||||
exp = self.exp_gen.gen(self.trace, selection)
|
||||
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
|
||||
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)
|
||||
elif isinstance(exp.sub_tasks[0], PipelineTask):
|
||||
exp = self.pipeline_coder.develop(exp)
|
||||
else:
|
||||
raise NotImplementedError(f"Unsupported component in DataScienceRDLoop: {exp.hypothesis.component}")
|
||||
exp.sub_tasks = []
|
||||
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():
|
||||
new_exp = self.runner.develop(exp)
|
||||
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"]
|
||||
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,
|
||||
)
|
||||
logger.log_object(feedback)
|
||||
return feedback
|
||||
|
||||
def record(self, prev_out: dict[str, Any]):
|
||||
# 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")
|
||||
|
||||
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)
|
||||
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
|
||||
from rdagent.scenarios.data_science.loop import DataScienceRDLoop
|
||||
|
||||
|
||||
def main(
|
||||
|
||||
@@ -171,7 +171,7 @@ class ExpGen(ABC):
|
||||
self.scen = scen
|
||||
|
||||
@abstractmethod
|
||||
def gen(self, trace: Trace, selection: tuple[int, ...] = (-1,)) -> Experiment:
|
||||
def gen(self, trace: Trace) -> Experiment:
|
||||
"""
|
||||
Generate the experiment based on the trace.
|
||||
|
||||
|
||||
@@ -0,0 +1,296 @@
|
||||
import shutil
|
||||
import subprocess
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional, Union
|
||||
|
||||
from rdagent.app.data_science.conf import DS_RD_SETTING
|
||||
from rdagent.components.coder.data_science.ensemble import EnsembleCoSTEER
|
||||
from rdagent.components.coder.data_science.ensemble.exp import EnsembleTask
|
||||
from rdagent.components.coder.data_science.feature import FeatureCoSTEER
|
||||
from rdagent.components.coder.data_science.feature.exp import FeatureTask
|
||||
from rdagent.components.coder.data_science.model import ModelCoSTEER
|
||||
from rdagent.components.coder.data_science.model.exp import ModelTask
|
||||
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
|
||||
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
|
||||
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, ExpGen
|
||||
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 DSTrace
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.idea_pool import DSKnowledgeBase
|
||||
|
||||
|
||||
class DataScienceRDLoop(RDLoop):
|
||||
# NOTE: we move the DataScienceRDLoop here to be easier to be imported
|
||||
skip_loop_error = (CoderError, RunnerError)
|
||||
|
||||
@staticmethod
|
||||
def _get_exp_gen(class_uri: str, scen: Scenario):
|
||||
"""
|
||||
Just for compatibility with the old version of the code.
|
||||
"""
|
||||
# TODO: remove me in the future. I don't have to be this complicated.
|
||||
# It is just for compatibility with the old version of the code and configuration.
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.proposal import (
|
||||
DSProposalV1ExpGen,
|
||||
DSProposalV2ExpGen,
|
||||
DSProposalV3ExpGen,
|
||||
)
|
||||
|
||||
if class_uri == "rdagent.scenarios.data_science.proposal.exp_gen.DSExpGen":
|
||||
if DS_RD_SETTING.proposal_version not in ["v1", "v2", "v3"]:
|
||||
return import_class(DS_RD_SETTING.proposal_version)(scen=scen)
|
||||
if DS_RD_SETTING.proposal_version == "v3":
|
||||
return DSProposalV3ExpGen(scen=scen)
|
||||
if DS_RD_SETTING.proposal_version == "v1":
|
||||
return DSProposalV1ExpGen(scen=scen)
|
||||
if DS_RD_SETTING.proposal_version == "v2":
|
||||
return DSProposalV2ExpGen(scen=scen)
|
||||
return import_class(class_uri)(scen)
|
||||
|
||||
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)
|
||||
|
||||
self.ckp_selector = import_class(PROP_SETTING.selector_name)()
|
||||
self.sota_exp_selector = import_class(PROP_SETTING.sota_exp_selector_name)()
|
||||
|
||||
self.exp_gen: ExpGen = self._get_exp_gen(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)
|
||||
|
||||
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]):
|
||||
|
||||
# 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
|
||||
selection = self.ckp_selector.get_selection(self.trace)
|
||||
# set the current selection for the trace
|
||||
self.trace.set_current_selection(selection)
|
||||
|
||||
exp = self.exp_gen.gen(self.trace)
|
||||
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
|
||||
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)
|
||||
elif isinstance(exp.sub_tasks[0], PipelineTask):
|
||||
exp = self.pipeline_coder.develop(exp)
|
||||
else:
|
||||
raise NotImplementedError(f"Unsupported component in DataScienceRDLoop: {exp.hypothesis.component}")
|
||||
exp.sub_tasks = []
|
||||
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():
|
||||
new_exp = self.runner.develop(exp)
|
||||
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"]
|
||||
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,
|
||||
)
|
||||
logger.log_object(feedback)
|
||||
return feedback
|
||||
|
||||
def record(self, prev_out: dict[str, Any]):
|
||||
# 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")
|
||||
|
||||
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)
|
||||
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
|
||||
@@ -1,47 +1,3 @@
|
||||
from rdagent.app.data_science.conf import DS_RD_SETTING
|
||||
from rdagent.core.proposal import ExpGen
|
||||
from rdagent.core.utils import import_class
|
||||
from rdagent.scenarios.data_science.experiment.experiment import DSExperiment
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.base import DSHypothesis, DSTrace
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.draft import DSDraftExpGen
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.proposal import (
|
||||
DSProposalV1ExpGen,
|
||||
DSProposalV2ExpGen,
|
||||
DSProposalV3ExpGen,
|
||||
)
|
||||
from rdagent.scenarios.data_science.scen import DataScienceScen
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.base import DSTrace
|
||||
|
||||
|
||||
class DSExpGen(ExpGen):
|
||||
"""
|
||||
Data Science Task Generator.
|
||||
This is a experiment router generator;
|
||||
"""
|
||||
|
||||
def __init__(self, scen: DataScienceScen) -> None:
|
||||
super().__init__(scen)
|
||||
|
||||
def gen(self, trace: DSTrace, selection: tuple[int, ...] = (-1,)) -> DSExperiment:
|
||||
|
||||
# set the current selection for the trace
|
||||
# handy design:dynamically change the "current selection" attribute of the trace, and we donot need to pass selection as an argument to other functions
|
||||
trace.set_current_selection(selection)
|
||||
|
||||
if DS_RD_SETTING.proposal_version not in ["v1", "v2", "v3"]:
|
||||
return import_class(DS_RD_SETTING.proposal_version)(scen=self.scen).gen(trace=trace)
|
||||
if DS_RD_SETTING.proposal_version == "v3":
|
||||
return DSProposalV3ExpGen(scen=self.scen).gen(trace=trace, pipeline=True)
|
||||
|
||||
if DS_RD_SETTING.coder_on_whole_pipeline:
|
||||
return DSProposalV2ExpGen(scen=self.scen).gen(trace=trace, pipeline=True)
|
||||
|
||||
next_missing_component = trace.next_incomplete_component()
|
||||
if next_missing_component is not None:
|
||||
return DSDraftExpGen(scen=self.scen).gen(
|
||||
component=next_missing_component,
|
||||
trace=trace,
|
||||
)
|
||||
if DS_RD_SETTING.proposal_version == "v1":
|
||||
return DSProposalV1ExpGen(scen=self.scen).gen(trace=trace)
|
||||
if DS_RD_SETTING.proposal_version == "v2":
|
||||
return DSProposalV2ExpGen(scen=self.scen).gen(trace=trace)
|
||||
__all__ = ["DSTrace"]
|
||||
|
||||
@@ -61,8 +61,6 @@ class DSTrace(Trace[DataScienceScen, KnowledgeBase]):
|
||||
|
||||
self.knowledge_base = knowledge_base
|
||||
|
||||
self.sub_trace_count: int = 0
|
||||
|
||||
self.current_selection: tuple[int, ...] = (-1,)
|
||||
|
||||
self.sota_exp_to_submit: DSExperiment | None = None # grab the global best exp to submit
|
||||
@@ -78,6 +76,10 @@ class DSTrace(Trace[DataScienceScen, KnowledgeBase]):
|
||||
def set_current_selection(self, selection: tuple[int, ...]) -> None:
|
||||
self.current_selection = selection
|
||||
|
||||
@property
|
||||
def sub_trace_count(self) -> int:
|
||||
return len(self.get_leaves())
|
||||
|
||||
def get_leaves(self) -> list[int, ...]:
|
||||
"""
|
||||
Get the indices of nodes (in hist) that have no children—i.e., "leaves" of current DAG.
|
||||
@@ -85,6 +87,10 @@ class DSTrace(Trace[DataScienceScen, KnowledgeBase]):
|
||||
tuple of ints: Indices of leaf nodes.
|
||||
- Leaves with lower index comes first.
|
||||
"""
|
||||
# BUG: potential BUG:
|
||||
# If we implement the most correct merging logic, merge 2 traces, will result in a single trace(2 traces currently).
|
||||
# So user may get unexpected results when he want to know ho many branches are created.
|
||||
|
||||
# Build a set of all parent indices found in dag_parent (skip empty tuples which represent roots)
|
||||
parent_indices = set(idx for parents in self.dag_parent for idx in parents)
|
||||
# All node indices
|
||||
@@ -153,12 +159,14 @@ class DSTrace(Trace[DataScienceScen, KnowledgeBase]):
|
||||
|
||||
def collect_all_ancestors(
|
||||
self,
|
||||
selection: tuple[int, ...] = (-1,),
|
||||
selection: tuple[int, ...] | None = None,
|
||||
) -> list[tuple[DSExperiment, ExperimentFeedback]]:
|
||||
"""
|
||||
Collect all ancestors of the given selection.
|
||||
The return list follows the order of [root->...->parent->current_node].
|
||||
"""
|
||||
if selection is None:
|
||||
selection = self.get_current_selection()
|
||||
|
||||
if len(self.dag_parent) == 0:
|
||||
return []
|
||||
|
||||
@@ -65,7 +65,6 @@ class LimitTimeCKPSelector(CheckpointSelector):
|
||||
current_time = datetime.now()
|
||||
|
||||
if len(trace.hist) == 0:
|
||||
trace.sub_trace_count = 0
|
||||
self.sub_trace_start_times[trace.sub_trace_count] = current_time
|
||||
logger.info(f"Starting initial sub-trace {trace.sub_trace_count} at {current_time}")
|
||||
return (-1,) # Continue with latest trial for new sub-trace
|
||||
@@ -90,12 +89,11 @@ class LimitTimeCKPSelector(CheckpointSelector):
|
||||
return (-1,)
|
||||
|
||||
# Time limit exceeded, start a new sub-trace
|
||||
trace.sub_trace_count += 1
|
||||
self.sub_trace_start_times[trace.sub_trace_count] = current_time
|
||||
self.sub_trace_start_times[trace.sub_trace_count + 1] = current_time
|
||||
logger.info(
|
||||
f"Elapsed time {elapsed_time} exceeds time limit {self.time_limit_pre_trace}, jump to a new sub-trace"
|
||||
)
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count}")
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count + 1}")
|
||||
return tuple() # Empty tuple signals starting a new sub-trace
|
||||
|
||||
|
||||
@@ -138,11 +136,10 @@ class SOTAJumpCKPSelector(CheckpointSelector):
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count}")
|
||||
return (-1,)
|
||||
|
||||
trace.sub_trace_count += 1
|
||||
logger.info(
|
||||
f"SOTA count {sota_count} is below threshold {self.SOTA_COUNT_THRESHOLD}, jump to a new sub-trace"
|
||||
)
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count}")
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count + 1}")
|
||||
return ()
|
||||
else:
|
||||
logger.info(
|
||||
@@ -201,7 +198,6 @@ class BackJumpCKPSelector(CheckpointSelector):
|
||||
|
||||
random_choice = random.random()
|
||||
if random_choice < 0.5:
|
||||
trace.sub_trace_count += 1
|
||||
logger.info(
|
||||
f"SOTA count {sota_count} is below threshold {self.SOTA_COUNT_THRESHOLD}, jump a new sub-trace"
|
||||
)
|
||||
@@ -227,11 +223,10 @@ class BackJumpCKPSelector(CheckpointSelector):
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count}")
|
||||
return (-1,)
|
||||
|
||||
trace.sub_trace_count += 1
|
||||
logger.info(
|
||||
f"SOTA count {sota_count} is below threshold {self.SOTA_COUNT_THRESHOLD}, jump a new sub-trace"
|
||||
)
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count}")
|
||||
logger.info(f"current sub-trace count: {trace.sub_trace_count + 1}")
|
||||
return () # reboot a new sub-trace
|
||||
|
||||
else:
|
||||
|
||||
@@ -8,13 +8,13 @@ from rdagent.core.proposal import ExpGen
|
||||
from rdagent.log import rdagent_logger as logger
|
||||
from rdagent.log.timer import RD_Agent_TIMER_wrapper, RDAgentTimer
|
||||
from rdagent.scenarios.data_science.experiment.experiment import DSExperiment
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen import DSExpGen
|
||||
from rdagent.scenarios.data_science.loop import DataScienceRDLoop
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.base import DSHypothesis, DSTrace
|
||||
from rdagent.utils.agent.tpl import T
|
||||
|
||||
|
||||
class MergeExpGen(ExpGen):
|
||||
def gen(self, trace: DSTrace, selection: tuple[int, ...] = (-1,)) -> DSExperiment:
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
# Ignore the selection argument and use all leaves instead.
|
||||
leaves: list[int] = trace.get_leaves()
|
||||
trace.set_current_selection((leaves[0],)) # override the current selection.
|
||||
@@ -87,9 +87,11 @@ class ExpGen2TraceAndMerge(ExpGen):
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.merge_exp_gen = MergeExpGen(self.scen)
|
||||
self.exp_gen = DSExpGen(self.scen)
|
||||
self.exp_gen = DataScienceRDLoop._get_exp_gen(
|
||||
"rdagent.scenarios.data_science.proposal.exp_gen.DSExpGen", self.scen
|
||||
)
|
||||
|
||||
def gen(self, trace: DSTrace, selection: tuple[int, ...] = (-1,)) -> DSExperiment:
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
timer: RDAgentTimer = RD_Agent_TIMER_wrapper.timer
|
||||
logger.info(f"Remain time: {timer.remain_time_duration}")
|
||||
|
||||
@@ -101,24 +103,23 @@ class ExpGen2TraceAndMerge(ExpGen):
|
||||
selection = (
|
||||
leaves[0],
|
||||
) # continue the first trace. This will result in the interleaving of two traces expansion.
|
||||
return self.exp_gen.gen(trace, selection)
|
||||
trace.set_current_selection(selection)
|
||||
return self.exp_gen.gen(trace)
|
||||
else:
|
||||
# disable reset in merging stage
|
||||
DS_RD_SETTING.coding_fail_reanalyze_threshold = 100000
|
||||
DS_RD_SETTING.consecutive_errors = 100000
|
||||
|
||||
leaves: list[int] = trace.get_leaves()
|
||||
if len(leaves) < 2:
|
||||
return self.exp_gen.gen(trace, selection)
|
||||
if trace.sub_trace_count < 2:
|
||||
return self.exp_gen.gen(trace)
|
||||
else:
|
||||
return self.merge_exp_gen.gen(trace, selection)
|
||||
return self.merge_exp_gen.gen(trace)
|
||||
|
||||
|
||||
class MergeExpGen_MultiTrace(ExpGen):
|
||||
def gen(self, trace: DSTrace, selection: tuple[int, ...] = (-1,)) -> DSExperiment:
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
# Ignore the selection argument and use all leaves instead.
|
||||
leaves: list[int] = trace.get_leaves()
|
||||
trace.set_current_selection(selection) #
|
||||
|
||||
# assuming merging the first and sencond trace.
|
||||
sota_exp_fb = trace.sota_experiment_fb(selection=(leaves[0],))
|
||||
@@ -195,11 +196,13 @@ class ExpGen2TraceAndMergeV2(ExpGen):
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.merge_exp_gen = MergeExpGen_MultiTrace(self.scen)
|
||||
self.exp_gen = DSExpGen(self.scen)
|
||||
self.exp_gen = DataScienceRDLoop._get_exp_gen(
|
||||
"rdagent.scenarios.data_science.proposal.exp_gen.DSExpGen", self.scen
|
||||
)
|
||||
self.MAX_TRACE_NUM = DS_RD_SETTING.max_trace_num # maximum number of traces to grow before merging
|
||||
self.flag_start_merge = False
|
||||
|
||||
def gen(self, trace: DSTrace, selection: tuple[int, ...] = (-1,)) -> DSExperiment:
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
timer: RDAgentTimer = RD_Agent_TIMER_wrapper.timer
|
||||
logger.info(f"Remain time: {timer.remain_time_duration}")
|
||||
|
||||
@@ -214,8 +217,7 @@ class ExpGen2TraceAndMergeV2(ExpGen):
|
||||
else:
|
||||
# set the knowledge base option back to False for the other traces
|
||||
DS_RD_SETTING.enable_knowledge_base = False
|
||||
|
||||
return self.exp_gen.gen(trace, selection)
|
||||
return self.exp_gen.gen(trace)
|
||||
|
||||
else:
|
||||
# disable reset in merging stage
|
||||
@@ -224,15 +226,14 @@ class ExpGen2TraceAndMergeV2(ExpGen):
|
||||
|
||||
leaves: list[int] = trace.get_leaves()
|
||||
if len(leaves) < 2:
|
||||
return self.exp_gen.gen(trace, selection=(-1,))
|
||||
trace.set_current_selection(selection=(-1,))
|
||||
return self.exp_gen.gen(trace)
|
||||
else:
|
||||
|
||||
if not self.flag_start_merge: # root node of the merge trace
|
||||
self.flag_start_merge = True
|
||||
selection = tuple()
|
||||
return self.merge_exp_gen.gen(trace, selection)
|
||||
trace.set_current_selection(tuple())
|
||||
return self.merge_exp_gen.gen(trace)
|
||||
else:
|
||||
# return self.merge_exp_gen.gen(trace, selection)
|
||||
return self.exp_gen.gen(
|
||||
trace, selection=(-1,)
|
||||
) # continue the last trace, to polish the merged solution
|
||||
# return self.merge_exp_gen.gen(trace)
|
||||
trace.set_current_selection(selection=(-1,))
|
||||
return self.exp_gen.gen(trace) # continue the last trace, to polish the merged solution
|
||||
|
||||
@@ -14,10 +14,12 @@ from rdagent.components.coder.data_science.pipeline.exp import PipelineTask
|
||||
from rdagent.components.coder.data_science.raw_data_loader.exp import DataLoaderTask
|
||||
from rdagent.components.coder.data_science.workflow.exp import WorkflowTask
|
||||
from rdagent.core.proposal import ExpGen
|
||||
from rdagent.core.scenario import Scenario
|
||||
from rdagent.log import rdagent_logger as logger
|
||||
from rdagent.oai.llm_utils import APIBackend, md5_hash
|
||||
from rdagent.scenarios.data_science.experiment.experiment import DSExperiment
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.base import DSHypothesis, DSTrace
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.draft import DSDraftExpGen
|
||||
from rdagent.scenarios.data_science.proposal.exp_gen.idea_pool import DSIdea
|
||||
from rdagent.utils.agent.tpl import T
|
||||
from rdagent.utils.repo.diff import generate_diff_from_dict
|
||||
@@ -259,8 +261,23 @@ COMPONENT_TASK_MAPPING = {
|
||||
}
|
||||
|
||||
|
||||
def draft_exp_in_decomposition(scen: Scenario, trace: DSTrace) -> None | DSDraftExpGen:
|
||||
next_missing_component = trace.next_incomplete_component()
|
||||
if next_missing_component is not None:
|
||||
return DSDraftExpGen(scen=scen).gen(
|
||||
component=next_missing_component,
|
||||
trace=trace,
|
||||
)
|
||||
else:
|
||||
return None
|
||||
|
||||
|
||||
class DSProposalV1ExpGen(ExpGen):
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
# Drafting Stage
|
||||
if draft_exp := draft_exp_in_decomposition(self.scen, trace):
|
||||
return draft_exp
|
||||
|
||||
# Guidelines:
|
||||
# System prompts: Shared condition you are facing
|
||||
# - scenario description: `scenario_desc`
|
||||
@@ -463,6 +480,36 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
)
|
||||
return json.loads(response)
|
||||
|
||||
def identify_problem(
|
||||
self, current_sub_trace, scenario_desc, sota_exp_desc, exp_feedback_list_desc, inject_diverse
|
||||
) -> Dict:
|
||||
sota_exp_num = sum(1 for _, fb in current_sub_trace if fb.decision)
|
||||
failed_exp_num = len(current_sub_trace) - sota_exp_num
|
||||
weighted_exp_num = (sota_exp_num * 3 + failed_exp_num * 2) // 2
|
||||
self.scen_prob_multiplier = max(0, 3 - weighted_exp_num // 4)
|
||||
|
||||
all_problems = {}
|
||||
if self.scen_prob_multiplier > 0:
|
||||
scen_problems = self.identify_scenario_problem(
|
||||
scenario_desc=scenario_desc,
|
||||
sota_exp_desc=sota_exp_desc,
|
||||
)
|
||||
for problem_name in scen_problems:
|
||||
scen_problems[problem_name]["label"] = "SCENARIO_PROBLEM"
|
||||
all_problems[problem_name] = scen_problems[problem_name]
|
||||
|
||||
if self.scen_prob_multiplier < 3:
|
||||
fb_problems = self.identify_feedback_problem(
|
||||
scenario_desc=scenario_desc,
|
||||
exp_feedback_list_desc=exp_feedback_list_desc,
|
||||
sota_exp_desc=sota_exp_desc,
|
||||
inject_diverse=inject_diverse,
|
||||
)
|
||||
for problem_name in fb_problems:
|
||||
fb_problems[problem_name]["label"] = "FEEDBACK_PROBLEM"
|
||||
all_problems[problem_name] = fb_problems[problem_name]
|
||||
return all_problems
|
||||
|
||||
@wait_retry(retry_n=5)
|
||||
def hypothesis_gen(
|
||||
self,
|
||||
@@ -477,7 +524,7 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
) -> Dict:
|
||||
problem_formatted_str = ""
|
||||
for problem_name, problem_dict in problems.items():
|
||||
problem_formatted_str += f"# Problem Name: {problem_name}\n"
|
||||
problem_formatted_str += f"Problem Name: {problem_name}\n"
|
||||
problem_formatted_str += f"- Problem Description: {problem_dict['problem']}\n"
|
||||
if "idea" in problem_dict:
|
||||
idea_formatted_str = DSIdea(problem_dict["idea"]).to_formatted_str()
|
||||
@@ -514,8 +561,10 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
self,
|
||||
hypothesis_dict: dict,
|
||||
problem_dict: dict,
|
||||
trace: DSTrace,
|
||||
) -> Tuple[str, DSHypothesis]:
|
||||
"""
|
||||
This function depends on the `identify_problem` function.
|
||||
"""
|
||||
weights = {
|
||||
"alignment_score": 0.2,
|
||||
"impact_score": 0.4,
|
||||
@@ -544,15 +593,15 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
scores_sorted = scores_sorted[:5] # Select top 5 hypotheses
|
||||
|
||||
# Increase the weight of the hypothesis that is inspired by the idea pool to 3x.
|
||||
# Linear decay the weight of the scenario problem from 3x to 1x.
|
||||
# Linear decay the weight of the scenario problem from 3x to 0x.
|
||||
index_to_pick_pool_list = []
|
||||
for j, problem_name in enumerate(scores_sorted.index):
|
||||
if hypothesis_dict[problem_name].get("inspired", False):
|
||||
index_to_pick_pool_list.extend([j] * 4)
|
||||
elif problem_dict.get(problem_name, {}).get("label", "") == "SCENARIO_PROBLEM":
|
||||
index_to_pick_pool_list.extend([j] * (3 - len(trace.hist) // 3))
|
||||
else:
|
||||
index_to_pick_pool_list.extend([j] * 2)
|
||||
if problem_dict.get(problem_name, {}).get("label", "") == "SCENARIO_PROBLEM":
|
||||
index_to_pick_pool_list.extend([j] * self.scen_prob_multiplier)
|
||||
else:
|
||||
index_to_pick_pool_list.extend([j] * (3 - self.scen_prob_multiplier))
|
||||
logger.info(f"index_to_pick_pool_list: {index_to_pick_pool_list}")
|
||||
|
||||
# Create a random but reproducible integer
|
||||
@@ -639,7 +688,10 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
exp.pending_tasks_list.append([workflow_task])
|
||||
return exp
|
||||
|
||||
def gen(self, trace: DSTrace, pipeline: bool = False) -> DSExperiment:
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
pipeline = DS_RD_SETTING.coder_on_whole_pipeline
|
||||
if not pipeline and (draft_exp := draft_exp_in_decomposition(self.scen, trace)):
|
||||
return draft_exp
|
||||
|
||||
if pipeline:
|
||||
component_desc = T("scenarios.data_science.share:component_description_in_pipeline").r()
|
||||
@@ -673,16 +725,6 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
pipeline=pipeline,
|
||||
)
|
||||
|
||||
if DS_RD_SETTING.enable_inject_diverse and len(trace.hist) > 0:
|
||||
if len(trace.current_selection) == 0:
|
||||
# start a new sub-trace, and inject diverse problems.
|
||||
inject_diverse = True
|
||||
logger.info("Start a new sub-trace, and inject diverse problems.")
|
||||
else:
|
||||
inject_diverse = False
|
||||
else:
|
||||
inject_diverse = False
|
||||
|
||||
if DS_RD_SETTING.enable_inject_diverse and len(trace.hist) > 0:
|
||||
if len(trace.current_selection) == 0:
|
||||
# start a new sub-trace, and inject diverse problems.
|
||||
@@ -694,27 +736,13 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
inject_diverse = False
|
||||
|
||||
# Step 1: Identify problems
|
||||
current_sub_trace = trace.collect_all_ancestors(selection=(-1,))
|
||||
all_problems = {}
|
||||
if len(current_sub_trace) >= 3:
|
||||
fb_problems = self.identify_feedback_problem(
|
||||
scenario_desc=scenario_desc,
|
||||
exp_feedback_list_desc=exp_feedback_list_desc,
|
||||
sota_exp_desc=sota_exp_desc,
|
||||
inject_diverse=inject_diverse,
|
||||
)
|
||||
for problem_name in fb_problems:
|
||||
fb_problems[problem_name]["label"] = "FEEDBACK_PROBLEM"
|
||||
all_problems[problem_name] = fb_problems[problem_name]
|
||||
|
||||
if len(current_sub_trace) < 9:
|
||||
scen_problems = self.identify_scenario_problem(
|
||||
scenario_desc=scenario_desc,
|
||||
sota_exp_desc=sota_exp_desc,
|
||||
)
|
||||
for problem_name in scen_problems:
|
||||
scen_problems[problem_name]["label"] = "SCENARIO_PROBLEM"
|
||||
all_problems[problem_name] = scen_problems[problem_name]
|
||||
all_problems = self.identify_problem(
|
||||
current_sub_trace=trace.collect_all_ancestors(),
|
||||
scenario_desc=scenario_desc,
|
||||
sota_exp_desc=sota_exp_desc,
|
||||
exp_feedback_list_desc=exp_feedback_list_desc,
|
||||
inject_diverse=inject_diverse,
|
||||
)
|
||||
|
||||
# Step 1.5: Sample ideas from idea pool
|
||||
if DS_RD_SETTING.enable_knowledge_base:
|
||||
@@ -757,7 +785,6 @@ class DSProposalV2ExpGen(ExpGen):
|
||||
pickled_problem_name, new_hypothesis = self.hypothesis_rank(
|
||||
hypothesis_dict=hypothesis_dict,
|
||||
problem_dict=all_problems,
|
||||
trace=trace,
|
||||
)
|
||||
# Step 3.5: Update knowledge base with the picked problem
|
||||
if DS_RD_SETTING.enable_knowledge_base:
|
||||
@@ -964,7 +991,11 @@ class DSProposalV3ExpGen(DSProposalV2ExpGen):
|
||||
)
|
||||
return result
|
||||
|
||||
def gen(self, trace: DSTrace, pipeline: bool = False) -> DSExperiment:
|
||||
def gen(self, trace: DSTrace) -> DSExperiment:
|
||||
pipeline = DS_RD_SETTING.coder_on_whole_pipeline
|
||||
if not pipeline and (draft_exp := draft_exp_in_decomposition(self.scen, trace)):
|
||||
return draft_exp
|
||||
|
||||
if pipeline:
|
||||
component_desc = T("scenarios.data_science.share:component_description_in_pipeline").r()
|
||||
else:
|
||||
|
||||
Reference in New Issue
Block a user