From c9346a9377d343ecc63be11bea6180ed5b34cdcf Mon Sep 17 00:00:00 2001 From: Xu Yang Date: Tue, 16 Jul 2024 10:33:53 +0800 Subject: [PATCH] first version of model runner and model feedback (#70) * Implemented model.py - Need to run within the RDAgent folder (relevant path) - Each time copy a template & insert code & run qlib & store result back to experiment * Create model.py * Create conf.yaml This is the sample conf.yaml to be copied each time. This has gone several times of iteration and is now working for both tabular and Time-Series data. * Create read_exp.py This is to read the results within Qlib * Create ReadMe.md * Update model.py * Create test_model.py A testing file that separates model code generation and running&feedback section. * move the template folder * help xisen finish the model runner * help xisen fix improve model feedback generation * delete debug file * rename readme.md --------- Co-authored-by: Xisen Wang <118058822+Xisen-Wang@users.noreply.github.com> --- rdagent/app/qlib_rd_loop/model.py | 2 +- .../coder/model_coder/CoSTEER/evaluators.py | 12 +- .../model_coder/CoSTEER/evolving_strategy.py | 3 +- rdagent/components/coder/model_coder/model.py | 4 +- rdagent/components/runner/__init__.py | 33 ++++++ .../qlib => components/runner}/conf.py | 4 +- rdagent/core/evolving_agent.py | 2 +- rdagent/core/experiment.py | 6 +- rdagent/core/proposal.py | 21 ++-- rdagent/core/task_generator.py | 6 +- rdagent/scenarios/qlib/prompts.yaml | 36 +++++- rdagent/scenarios/qlib/task_generator/data.py | 53 ++------- .../task_generator/env_factor/read_exp_res.py | 14 +-- .../scenarios/qlib/task_generator/feedback.py | 94 +++++---------- .../scenarios/qlib/task_generator/model.py | 73 +++++++++++- .../task_generator/model_template/README.md | 3 + .../task_generator/model_template/conf.yaml | 109 ++++++++++++++++++ .../model_template/read_exp_res.py | 45 ++++++++ rdagent/utils/env.py | 3 + 19 files changed, 369 insertions(+), 154 deletions(-) create mode 100644 rdagent/components/runner/__init__.py rename rdagent/{scenarios/qlib => components/runner}/conf.py (84%) create mode 100644 rdagent/scenarios/qlib/task_generator/model_template/README.md create mode 100644 rdagent/scenarios/qlib/task_generator/model_template/conf.yaml create mode 100644 rdagent/scenarios/qlib/task_generator/model_template/read_exp_res.py diff --git a/rdagent/app/qlib_rd_loop/model.py b/rdagent/app/qlib_rd_loop/model.py index cd3bb605..e238d6b9 100644 --- a/rdagent/app/qlib_rd_loop/model.py +++ b/rdagent/app/qlib_rd_loop/model.py @@ -25,7 +25,7 @@ hypothesis2experiment: Hypothesis2Experiment = import_class(PROP_SETTING.qlib_mo qlib_model_coder: TaskGenerator = import_class(PROP_SETTING.qlib_model_coder)(scen) qlib_model_runner: TaskGenerator = import_class(PROP_SETTING.qlib_model_runner)(scen) -qlib_model_summarizer: HypothesisExperiment2Feedback = import_class(PROP_SETTING.qlib_model_summarizer)() +qlib_model_summarizer: HypothesisExperiment2Feedback = import_class(PROP_SETTING.qlib_model_summarizer)(scen) trace = Trace(scen=scen) for _ in range(PROP_SETTING.evolving_n): diff --git a/rdagent/components/coder/model_coder/CoSTEER/evaluators.py b/rdagent/components/coder/model_coder/CoSTEER/evaluators.py index 168251c5..c3e5f468 100644 --- a/rdagent/components/coder/model_coder/CoSTEER/evaluators.py +++ b/rdagent/components/coder/model_coder/CoSTEER/evaluators.py @@ -241,12 +241,12 @@ class ModelCoderEvaluator(Evaluator): ) assert isinstance(target_task, ModelTask) - batch_size, num_features, num_timesteps = ( - random.randint(6, 10), - random.randint(6, 10), - random.randint(6, 10), - ) - input_value, param_init_value = random.random(), random.random() + # NOTE: Use fixed input to test the model to avoid randomness + batch_size = 8 + num_features = 30 + num_timesteps = 40 + input_value = 0.4 + param_init_value = 0.6 assert isinstance(implementation, ModelImplementation) model_execution_feedback, gen_tensor = implementation.execute( diff --git a/rdagent/components/coder/model_coder/CoSTEER/evolving_strategy.py b/rdagent/components/coder/model_coder/CoSTEER/evolving_strategy.py index 71e9e5e7..07a6162c 100644 --- a/rdagent/components/coder/model_coder/CoSTEER/evolving_strategy.py +++ b/rdagent/components/coder/model_coder/CoSTEER/evolving_strategy.py @@ -86,13 +86,12 @@ class ModelCoderEvolvingStrategy(EvolvingStrategy): queried_similar_successful_knowledge_to_render = queried_similar_successful_knowledge_to_render[1:] code = json.loads( - APIBackend(use_chat_cache=True).build_messages_and_create_chat_completion( + APIBackend(use_chat_cache=False).build_messages_and_create_chat_completion( user_prompt=user_prompt, system_prompt=system_prompt, json_mode=True, ), )["code"] - # ast.parse(code) model_implementation = ModelImplementation( target_task, ) diff --git a/rdagent/components/coder/model_coder/model.py b/rdagent/components/coder/model_coder/model.py index 7ddf634b..e80272d2 100644 --- a/rdagent/components/coder/model_coder/model.py +++ b/rdagent/components/coder/model_coder/model.py @@ -83,7 +83,9 @@ class ModelImplementation(FBImplementation): try: if MODEL_IMPL_SETTINGS.enable_execution_cache: # NOTE: cache the result for the same code - target_file_name = md5_hash(self.code_dict["model.py"]) + target_file_name = md5_hash( + f"{batch_size}_{num_features}_{num_timesteps}_{input_value}_{param_init_value}_{self.code_dict['model.py']}" + ) cache_file_path = Path(MODEL_IMPL_SETTINGS.model_cache_location) / f"{target_file_name}.pkl" Path(MODEL_IMPL_SETTINGS.model_cache_location).mkdir(exist_ok=True, parents=True) if cache_file_path.exists(): diff --git a/rdagent/components/runner/__init__.py b/rdagent/components/runner/__init__.py new file mode 100644 index 00000000..d331f18b --- /dev/null +++ b/rdagent/components/runner/__init__.py @@ -0,0 +1,33 @@ +import pickle +from pathlib import Path +from typing import Tuple + +from rdagent.components.runner.conf import RUNNER_SETTINGS +from rdagent.core.experiment import ASpecificExp, Experiment +from rdagent.core.task_generator import TaskGenerator +from rdagent.oai.llm_utils import md5_hash + + +class CachedRunner(TaskGenerator[ASpecificExp]): + def get_cache_key(self, exp: Experiment) -> str: + all_tasks = [] + for based_exp in exp.based_experiments: + all_tasks.extend(based_exp.sub_tasks) + all_tasks.extend(exp.sub_tasks) + task_info_list = [task.get_task_information() for task in all_tasks] + task_info_str = "\n".join(task_info_list) + return md5_hash(task_info_str) + + def get_cache_result(self, exp: Experiment) -> Tuple[bool, object]: + task_info_key = self.get_cache_key(exp) + Path(RUNNER_SETTINGS.runner_cache_path).mkdir(parents=True, exist_ok=True) + cache_path = Path(RUNNER_SETTINGS.runner_cache_path) / f"{task_info_key}.pkl" + if cache_path.exists(): + return True, pickle.load(open(cache_path, "rb")) + else: + return False, None + + def dump_cache_result(self, exp: Experiment, result: object): + task_info_key = self.get_cache_key(exp) + cache_path = Path(RUNNER_SETTINGS.runner_cache_path) / f"{task_info_key}.pkl" + pickle.dump(result, open(cache_path, "wb")) diff --git a/rdagent/scenarios/qlib/conf.py b/rdagent/components/runner/conf.py similarity index 84% rename from rdagent/scenarios/qlib/conf.py rename to rdagent/components/runner/conf.py index e4fbc363..b7575d57 100644 --- a/rdagent/scenarios/qlib/conf.py +++ b/rdagent/components/runner/conf.py @@ -11,9 +11,9 @@ load_dotenv(verbose=True, override=True) from pydantic_settings import BaseSettings -class QlibRDAgentSettings(BaseSettings): +class RunnerSettings(BaseSettings): runner_cache_result: bool = True # whether to cache the result of the docker execution runner_cache_path: str = str(Path.cwd() / "runner_cache/") # the path to store the cache -Qlib_RD_AGENT_SETTINGS = QlibRDAgentSettings() +RUNNER_SETTINGS = RunnerSettings() diff --git a/rdagent/core/evolving_agent.py b/rdagent/core/evolving_agent.py index 855d7e40..42c1e071 100644 --- a/rdagent/core/evolving_agent.py +++ b/rdagent/core/evolving_agent.py @@ -32,7 +32,7 @@ class RAGEvoAgent(EvoAgent): with_feedback: bool = True, knowledge_self_gen: bool = False, ) -> EvolvableSubjects: - for _ in tqdm(range(self.max_loop), "Implementing factors"): + for _ in tqdm(range(self.max_loop), "Implementing"): # 1. knowledge self-evolving if knowledge_self_gen and self.rag is not None: self.rag.generate_knowledge(self.evolving_trace) diff --git a/rdagent/core/experiment.py b/rdagent/core/experiment.py index a71545dc..c0c661ff 100644 --- a/rdagent/core/experiment.py +++ b/rdagent/core/experiment.py @@ -124,16 +124,16 @@ class Experiment(ABC, Generic[ASpecificTask, ASpecificImp]): The experiment is a sequence of tasks and the implementations of the tasks after generated by the TaskGenerator. """ - result_ws: Optional[FBImplementation] - def __init__(self, sub_tasks: Sequence[ASpecificTask]) -> None: self.sub_tasks = sub_tasks self.sub_implementations: Sequence[ASpecificImp] = [None for _ in self.sub_tasks] self.based_experiments: Sequence[Experiment] = [] self.result: object = None # The result of the experiment, can be different types in different scenarios. - self.result_ws = None + self.exp_ws: ASpecificImp = None +ASpecificExp = TypeVar("ASpecificExp", bound=Experiment) + TaskOrExperiment = TypeVar("TaskOrExperiment", Task, Experiment) diff --git a/rdagent/core/proposal.py b/rdagent/core/proposal.py index 9b5e1bad..7dde0e20 100644 --- a/rdagent/core/proposal.py +++ b/rdagent/core/proposal.py @@ -6,7 +6,7 @@ from abc import ABC, abstractmethod from typing import Any, Dict, Generic, List, Tuple, TypeVar from rdagent.core.evaluation import Feedback -from rdagent.core.experiment import ASpecificTask, Experiment +from rdagent.core.experiment import ASpecificExp, ASpecificTask, Experiment from rdagent.core.scenario import Scenario # class data_ana: XXX @@ -23,7 +23,7 @@ class Hypothesis: def __init__(self, hypothesis: str, reason: str) -> None: self.hypothesis: str = hypothesis self.reason: str = reason - + def __str__(self) -> str: return f"""Hypothesis: {self.hypothesis} Reason: {self.reason}""" @@ -54,15 +54,15 @@ class Trace(Generic[ASpecificScen]): self.scen: ASpecificScen = scen self.hist: list[Tuple[Hypothesis, Experiment, HypothesisFeedback]] = [] - def get_last_experiment_info(self) -> Tuple[Hypothesis, ASpecificTask, Any]: + def get_SOTA_hypothesis_and_experiment(self) -> Tuple[Hypothesis, Experiment]: """Access the last experiment result, sub-task, and the corresponding hypothesis.""" # TODO: The return value does not align with the signature. - if not self.hist: - return None - last_hypothesis, last_experiment, _ = self.hist[-1] - last_task = last_experiment.sub_tasks[-1] - last_result = last_experiment.result - return last_hypothesis, last_task, last_result + for hypothesis, experiment, feedback in self.hist[::-1]: + if feedback.decision: + return hypothesis, experiment + + return None, None + class HypothesisGen: def __init__(self, scen: Scenario): @@ -81,9 +81,6 @@ class HypothesisGen: """ -ASpecificExp = TypeVar("ASpecificExp", bound=Experiment) - - class Hypothesis2Experiment(ABC, Generic[ASpecificExp]): """ [Abstract description => concrete description] => Code implement diff --git a/rdagent/core/task_generator.py b/rdagent/core/task_generator.py index 17a69d39..27c91ac4 100644 --- a/rdagent/core/task_generator.py +++ b/rdagent/core/task_generator.py @@ -1,11 +1,9 @@ from abc import ABC, abstractmethod -from typing import Generic, List, Sequence, TypeVar +from typing import Generic, List -from rdagent.core.experiment import Experiment +from rdagent.core.experiment import ASpecificExp from rdagent.core.scenario import Scenario -ASpecificExp = TypeVar("ASpecificExp", bound=Experiment) - class TaskGenerator(ABC, Generic[ASpecificExp]): def __init__(self, scen: Scenario) -> None: diff --git a/rdagent/scenarios/qlib/prompts.yaml b/rdagent/scenarios/qlib/prompts.yaml index f917b220..22830aaa 100644 --- a/rdagent/scenarios/qlib/prompts.yaml +++ b/rdagent/scenarios/qlib/prompts.yaml @@ -49,7 +49,7 @@ model_experiment_output_format: |- # Don't add ellipsis (...) or any filler text that might cause JSON parsing errors here! } -data_feedback_generation: +factor_feedback_generation: system: |- You are a professional result analysis assistant on data driven R&D. The task is described in the following scenario: @@ -80,3 +80,37 @@ data_feedback_generation: 3. Demonstrate positive or negative effects when compared to Alpha158. Provide detailed feedback and recommend whether to replace the best result if the new factor proves superior. + +model_feedback_generation: + system: |- + You are a professional result analysis assistant. You will receive a result and a hypothesis. + Your task is to provide feedback on how well the result supports or refutes the hypothesis by judging from the observation of performance increase or decrease. + Please provide detailed and constructive feedback. + Example JSON Structure for Result Analysis: + { + "Observations": "Your overall observations here", + "Feedback for Hypothesis": "Observations related to the hypothesis", + "New Hypothesis": "Put your new hypothesis here.", + "Reasoning": "Provide reasoning for the hypothesis here.", + "Decision": True or False, + } + user: |- + We are in an experiment of finding hypothesis and validating or rejecting them so that in the end we have a powerful model generated. + Here are the context: {{context}}. + + {% if last_hypothesis %} + Last Round Information: + Hypothesis: {{last_hypothesis.hypothesis}} + Task: {{last_task}} + Result: {{last_result}} + {% else %} + This is the first round. No previous information available. + {% endif %} + + Now let's come to this round. You will receive the result and you will evaluate if the performance increases or decreases. + Hypothesis: {{hypothesis.hypothesis}} + Relevant Reasoning: {{hypothesis.reason}} + Result: {{exp.result}} + + Compare and observe. Which result has a better return and lower risk? If the performance increases, the hypothesis should be considered positive (working). + Hence, with the hypotheses, relevant reasoning, and results in mind (comparison), provide detailed and constructive feedback and suggest a new hypothesis. \ No newline at end of file diff --git a/rdagent/scenarios/qlib/task_generator/data.py b/rdagent/scenarios/qlib/task_generator/data.py index bc44f203..2203b43d 100644 --- a/rdagent/scenarios/qlib/task_generator/data.py +++ b/rdagent/scenarios/qlib/task_generator/data.py @@ -5,10 +5,11 @@ from typing import List, Tuple import pandas as pd +from rdagent.components.runner import CachedRunner +from rdagent.components.runner.conf import RUNNER_SETTINGS from rdagent.core.log import RDAgentLog from rdagent.core.task_generator import TaskGenerator from rdagent.oai.llm_utils import md5_hash -from rdagent.scenarios.qlib.conf import Qlib_RD_AGENT_SETTINGS from rdagent.scenarios.qlib.experiment.factor_experiment import QlibFactorExperiment from rdagent.utils.env import QTDockerEnv @@ -30,7 +31,7 @@ logger = RDAgentLog() # TODO: supporting multiprocessing and keep previous results -class QlibFactorRunner(TaskGenerator[QlibFactorExperiment]): +class QlibFactorRunner(CachedRunner[QlibFactorExperiment]): """ Docker run Everything in a folder @@ -40,29 +41,6 @@ class QlibFactorRunner(TaskGenerator[QlibFactorExperiment]): - results in `mlflow` """ - def get_cache_key(self, exp: QlibFactorExperiment) -> str: - all_tasks = [] - for based_exp in exp.based_experiments: - all_tasks.extend(based_exp.sub_tasks) - all_tasks.extend(exp.sub_tasks) - task_info_list = [task.get_task_information() for task in all_tasks] - task_info_str = "\n".join(task_info_list) - return md5_hash(task_info_str) - - def get_cache_result(self, exp: QlibFactorExperiment) -> Tuple[bool, object]: - task_info_key = self.get_cache_key(exp) - Path(Qlib_RD_AGENT_SETTINGS.runner_cache_path).mkdir(parents=True, exist_ok=True) - cache_path = Path(Qlib_RD_AGENT_SETTINGS.runner_cache_path) / f"{task_info_key}.pkl" - if cache_path.exists(): - return True, pickle.load(open(cache_path, "rb")) - else: - return False, None - - def dump_cache_result(self, exp: QlibFactorExperiment, result: object): - task_info_key = self.get_cache_key(exp) - cache_path = Path(Qlib_RD_AGENT_SETTINGS.runner_cache_path) / f"{task_info_key}.pkl" - pickle.dump(result, open(cache_path, "wb")) - def generate(self, exp: QlibFactorExperiment) -> QlibFactorExperiment: """ Generate the experiment by processing and combining factor data, @@ -71,7 +49,7 @@ class QlibFactorRunner(TaskGenerator[QlibFactorExperiment]): if exp.based_experiments and exp.based_experiments[-1].result is None: exp.based_experiments[-1] = self.generate(exp.based_experiments[-1]) - if Qlib_RD_AGENT_SETTINGS.runner_cache_result: + if RUNNER_SETTINGS.runner_cache_result: cache_hit, result = self.get_cache_result(exp) if cache_hit: exp.result = result @@ -116,30 +94,19 @@ class QlibFactorRunner(TaskGenerator[QlibFactorExperiment]): execute_log = qtde.run(local_path=str(DIRNAME / "env_factor"), entry="python read_exp_res.py") - pkl_path = DIRNAME / "env_factor/qlib_res.pkl" + csv_path = DIRNAME / "env_factor/qlib_res.csv" - if not pkl_path.exists(): - logger.error(f"File {pkl_path} does not exist.") + if not csv_path.exists(): + logger.error(f"File {csv_path} does not exist.") return None - with open(pkl_path, "rb") as f: - result = pickle.load(f) + result = pd.read_csv(csv_path, index_col=0).iloc[:, 0] exp.result = result - if Qlib_RD_AGENT_SETTINGS.runner_cache_result: + if RUNNER_SETTINGS.runner_cache_result: self.dump_cache_result(exp, result) - # Check if the result is valid and is a DataFrame - if isinstance(result, pd.DataFrame): - if not result.empty: - logger.info("Successfully retrieved experiment result.") - return exp - else: - logger.error("Result DataFrame is empty.") - return None - else: - logger.error("Data format error.") - return None + return exp def process_factor_data(self, exp_or_list: List[QlibFactorExperiment] | QlibFactorExperiment) -> pd.DataFrame: """ diff --git a/rdagent/scenarios/qlib/task_generator/env_factor/read_exp_res.py b/rdagent/scenarios/qlib/task_generator/env_factor/read_exp_res.py index 6e05ec5e..19930248 100644 --- a/rdagent/scenarios/qlib/task_generator/env_factor/read_exp_res.py +++ b/rdagent/scenarios/qlib/task_generator/env_factor/read_exp_res.py @@ -1,4 +1,3 @@ -import os import pickle from pathlib import Path @@ -22,7 +21,6 @@ experiments = R.list_experiments() experiment_name = None latest_recorder = None for experiment in experiments: - # print(f"Experiment: {experiment}") recorders = R.list_recorders(experiment_name=experiment) for recorder_id in recorders: if recorder_id is not None: @@ -39,15 +37,9 @@ else: print(f"Latest recorder: {latest_recorder}") # Load the specified file from the latest recorder - file_path = "portfolio_analysis/port_analysis_1day.pkl" - indicator_analysis_df = latest_recorder.load_object(file_path) + metrics = pd.Series(latest_recorder.list_metrics()) - # Optionally convert to DataFrame if not already in DataFrame format - if not isinstance(indicator_analysis_df, pd.DataFrame): - indicator_analysis_df = pd.DataFrame(indicator_analysis_df) - - output_path = os.path.join(str(Path(__file__).resolve().parent), "qlib_res.pkl") - with open(output_path, "wb") as f: - pickle.dump(indicator_analysis_df, f) + output_path = Path(__file__).resolve().parent / "qlib_res.csv" + metrics.to_csv(output_path) print(f"Output has been saved to {output_path}") diff --git a/rdagent/scenarios/qlib/task_generator/feedback.py b/rdagent/scenarios/qlib/task_generator/feedback.py index 4629f2cc..8ce1e08e 100644 --- a/rdagent/scenarios/qlib/task_generator/feedback.py +++ b/rdagent/scenarios/qlib/task_generator/feedback.py @@ -44,14 +44,14 @@ class QlibFactorHypothesisExperiment2Feedback(HypothesisExperiment2Feedback): # Generate the system prompt sys_prompt = ( Environment(undefined=StrictUndefined) - .from_string(feedback_prompts["data_feedback_generation"]["system"]) + .from_string(feedback_prompts["factor_feedback_generation"]["system"]) .render(scenario=self.scen.get_scenario_all_desc()) ) # Generate the user prompt usr_prompt = ( Environment(undefined=StrictUndefined) - .from_string(feedback_prompts["data_feedback_generation"]["user"]) + .from_string(feedback_prompts["factor_feedback_generation"]["user"]) .render( hypothesis_text=hypothesis_text, task_details=tasks_factors, @@ -108,70 +108,38 @@ class QlibModelHypothesisExperiment2Feedback(HypothesisExperiment2Feedback): """ # Define the system prompt for hypothesis feedback - sys_prompt_hypothesis = ( - "You are a professional result analysis assistant. You will receive a result and a hypothesis. " - "Your task is to provide feedback on how well the result supports or refutes the hypothesis by judging from the observation of performance increase or decrease. " - "Please provide detailed and constructive feedback. " - "Example JSON Structure for Result Analysis: " - '{"Observations": "Your overall observations here", "Feedback for Hypothesis": "Observations related to the hypothesis", ' - '"New Hypothesis": "Put your new hypothesis here.", "Reasoning": "Provide reasoning for the hypothesis here.", ' - '"Decision": "True or False"}' - ) + system_prompt = feedback_prompts["model_feedback_generation"]["system"] # Define the user prompt for hypothesis feedback context = trace.scen - last_experiment_info = trace.get_last_experiment_info() + SOTA_hypothesis, SOTA_experiment = trace.get_SOTA_hypothesis_and_experiment() - if last_experiment_info: - last_hypothesis, last_task, last_result = last_experiment_info - last_info_str = f"Last Round Information:\nHypothesis: {last_hypothesis.hypothesis}\nTask: {last_task}\nResult: {last_result}\n" - else: - last_info_str = "This is the first round. No previous information available." - - usr_prompt_hypothesis = f""" - We are in an experiment of finding hypothesis and validating or rejecting them so that in the end we have a powerful model generated. - Here are the context: {context}. - {last_info_str} - - Now let's come to this round. You will receive the result and you will evaluate if the performance increases or decreases. - Hypothesis: {hypothesis.hypothesis}\n - Relevant Reasoning: {hypothesis.reason}\n - Result: {exp.result}\n - - Compare and observe. Which result has a better return and lower risk? If the performance increases, the hypothesis should be considered positive (working). - Hence, with the hypotheses, relevant reasoning, and results in mind (comparison), provide detailed and constructive feedback and suggest a new hypothesis. - """ - - try: - # Call the APIBackend to generate the response for hypothesis feedback - response_hypothesis = APIBackend().build_messages_and_create_chat_completion( - user_prompt=usr_prompt_hypothesis, - system_prompt=sys_prompt_hypothesis, - json_mode=True, + user_prompt = ( + Environment(undefined=StrictUndefined) + .from_string(feedback_prompts["model_feedback_generation"]["user"]) + .render( + context=context, + last_hypothesis=SOTA_hypothesis, + last_task=SOTA_experiment.sub_tasks[0].get_task_information() if SOTA_hypothesis else None, + last_result=SOTA_experiment.result if SOTA_hypothesis else None, + hypothesis=hypothesis, + exp=exp, ) - - # Parse the JSON response to extract the feedback - response_json_hypothesis = json.loads(response_hypothesis) - hypothesis_feedback = HypothesisFeedback( - observations=response_json_hypothesis.get("Observations", "No observations provided"), - hypothesis_evaluation=response_json_hypothesis.get("Feedback for Hypothesis", "No feedback provided"), - new_hypothesis=response_json_hypothesis.get("New Hypothesis", "No new hypothesis provided"), - reason=response_json_hypothesis.get("Reasoning", "No reasoning provided"), - decision=response_json_hypothesis.get("Decision", "false").lower() == "true", - ) - - return hypothesis_feedback - - except json.JSONDecodeError as e: - # TODO: (Xiao) I think raising a specific type of ERROR to make caller know sth bad has happened would be more reasonable - print("Error parsing JSON response from LLM for hypothesis feedback:", e) - except Exception as e: - print("An unexpected error occurred while generating hypothesis feedback:", e) - - return HypothesisFeedback( - observations="No observations", - hypothesis_evaluation="No feedback", - new_hypothesis="No new hypothesis", - reason="No reasoning", - decision=False, + ) + + # Call the APIBackend to generate the response for hypothesis feedback + response_hypothesis = APIBackend().build_messages_and_create_chat_completion( + user_prompt=user_prompt, + system_prompt=system_prompt, + json_mode=True, + ) + + # Parse the JSON response to extract the feedback + response_json_hypothesis = json.loads(response_hypothesis) + return HypothesisFeedback( + observations=response_json_hypothesis.get("Observations", "No observations provided"), + hypothesis_evaluation=response_json_hypothesis.get("Feedback for Hypothesis", "No feedback provided"), + new_hypothesis=response_json_hypothesis.get("New Hypothesis", "No new hypothesis provided"), + reason=response_json_hypothesis.get("Reasoning", "No reasoning provided"), + decision=response_json_hypothesis.get("Decision", "false").lower() == "true", ) diff --git a/rdagent/scenarios/qlib/task_generator/model.py b/rdagent/scenarios/qlib/task_generator/model.py index c7f336e4..47586924 100644 --- a/rdagent/scenarios/qlib/task_generator/model.py +++ b/rdagent/scenarios/qlib/task_generator/model.py @@ -1,8 +1,21 @@ -from rdagent.components.coder.model_coder.model import ModelImplementation +import shutil +import uuid +from pathlib import Path + +import pandas as pd + +from rdagent.components.coder.model_coder.model import ( + ModelExperiment, + ModelImplementation, +) +from rdagent.components.runner import CachedRunner +from rdagent.components.runner.conf import RUNNER_SETTINGS +from rdagent.core.log import RDAgentLog from rdagent.core.task_generator import TaskGenerator +from rdagent.utils.env import QTDockerEnv -class QlibModelRunner(TaskGenerator[ModelImplementation]): +class QlibModelRunner(CachedRunner[ModelImplementation]): """ Docker run Everything in a folder @@ -15,5 +28,57 @@ class QlibModelRunner(TaskGenerator[ModelImplementation]): - let LLM modify model.py """ - def generate(self, exp: ModelImplementation) -> ModelImplementation: - return exp # TODO IMPLEMENT THIS + def generate(self, exp: ModelExperiment) -> ModelExperiment: + + if RUNNER_SETTINGS.runner_cache_result: + cache_hit, result = self.get_cache_result(exp) + if cache_hit: + exp.result = result + return exp + TEMPLATE_PATH = Path(__file__).parent / "model_template" # Can be updated + + # To set the experiment level workspace and prepare the workspaces use the first task as the target task + exp.exp_ws = ModelImplementation(target_task=exp.sub_tasks[0]) + exp.exp_ws.prepare() + + # to copy_template_to_workspace + for file_path in TEMPLATE_PATH.iterdir(): + shutil.copyfile(file_path, exp.exp_ws.workspace_path / file_path.name) + + # to replace & inject code + exp.exp_ws.inject_code(**{"model.py": exp.sub_implementations[0].code_dict["model.py"]}) + + env_to_use = {} + + if exp.sub_tasks[0].model_type == "TimeSeries": + env_to_use = {"dataset_cls": "TSDatasetH", "step_len": 20, "num_timesteps": 20} + elif exp.sub_tasks[0].model_type == "Tabular": + env_to_use = {"dataset_cls": "DatasetH"} + # to execute + qtde = QTDockerEnv() + + # Preparing the Docker environment + qtde.prepare() + + # Run the Docker container with the specified entry + + execute_log = qtde.run( + local_path=exp.exp_ws.workspace_path, entry="qrun conf.yaml", env={"PYTHONPATH": "./", **env_to_use} + ) + + # Run the experiment analysis code + execute_log = qtde.run(local_path=exp.exp_ws.workspace_path, entry="python read_exp_res.py") + + csv_path = exp.exp_ws.workspace_path / "qlib_res.csv" + + if not csv_path.exists(): + RDAgentLog().error(f"File {csv_path} does not exist.") + return None + + result = pd.read_csv(csv_path, index_col=0).iloc[:,0] + + exp.result = result + if RUNNER_SETTINGS.runner_cache_result: + self.dump_cache_result(exp, result) + + return exp diff --git a/rdagent/scenarios/qlib/task_generator/model_template/README.md b/rdagent/scenarios/qlib/task_generator/model_template/README.md new file mode 100644 index 00000000..fb7c22c9 --- /dev/null +++ b/rdagent/scenarios/qlib/task_generator/model_template/README.md @@ -0,0 +1,3 @@ +## This folder is a template to be copied from for each model implementation & running process. + +Components: Dummy model.py, versatile conf.yaml, and a result reader. diff --git a/rdagent/scenarios/qlib/task_generator/model_template/conf.yaml b/rdagent/scenarios/qlib/task_generator/model_template/conf.yaml new file mode 100644 index 00000000..d735188e --- /dev/null +++ b/rdagent/scenarios/qlib/task_generator/model_template/conf.yaml @@ -0,0 +1,109 @@ +qlib_init: + provider_uri: "~/.qlib/qlib_data/cn_data" + region: cn +market: &market csi300 +benchmark: &benchmark SH000300 +data_handler_config: &data_handler_config + start_time: 2008-01-01 + end_time: 2020-08-01 + fit_start_time: 2008-01-01 + fit_end_time: 2014-12-31 + instruments: *market + infer_processors: + - class: FilterCol + kwargs: + fields_group: feature + col_list: ["RESI5", "WVMA5", "RSQR5", "KLEN", "RSQR10", "CORR5", "CORD5", "CORR10", + "ROC60", "RESI10", "VSTD5", "RSQR60", "CORR60", "WVMA60", "STD5", + "RSQR20", "CORD60", "CORD10", "CORR20", "KLOW" + ] + - class: RobustZScoreNorm + kwargs: + fields_group: feature + clip_outlier: true + - class: Fillna + kwargs: + fields_group: feature + learn_processors: + - class: DropnaLabel + - class: CSRankNorm + kwargs: + fields_group: label + label: ["Ref($close, -2) / Ref($close, -1) - 1"] + +port_analysis_config: &port_analysis_config + strategy: + class: TopkDropoutStrategy + module_path: qlib.contrib.strategy + kwargs: + signal: + topk: 50 + n_drop: 5 + backtest: + start_time: 2017-01-01 + end_time: 2020-08-01 + account: 100000000 + benchmark: *benchmark + exchange_kwargs: + limit_threshold: 0.095 + deal_price: close + open_cost: 0.0005 + close_cost: 0.0015 + min_cost: 5 +task: + model: + class: GeneralPTNN + module_path: qlib.contrib.model.pytorch_general_nn + kwargs: + n_epochs: 1 + lr: 2e-3 + early_stop: 10 + batch_size: 800 + metric: loss + loss: mse + n_jobs: 20 + GPU: 0 + # loss: mse + # lr: 0.002 + # optimizer: adam + # batch_size: 8192 + # GPU: 0 + # weight_decay: 0.0002 + # pt_model_uri: "qlib.contrib.model.pytorch_nn.Net" + # pt_model_uri: "env_tpl.model.Net" + # pt_model_uri: "./model.py:Net" + pt_model_uri: "model.model_cls" + pt_model_kwargs: { + "num_features": 20, + {% if num_timesteps %}num_timesteps: {{ num_timesteps }}{% endif %} + } + # input_dim: 20 + # How should I use jinja to put step len here conditionally + dataset: + class: {{ dataset_cls | default("DatasetH") }} + module_path: qlib.data.dataset + kwargs: + handler: + class: Alpha158 + module_path: qlib.contrib.data.handler + kwargs: *data_handler_config + segments: + train: [2008-01-01, 2014-12-31] + valid: [2015-01-01, 2016-12-31] + test: [2017-01-01, 2020-08-01] + {% if step_len %}step_len: {{ step_len }}{% endif %} + record: + - class: SignalRecord + module_path: qlib.workflow.record_temp + kwargs: + model: + dataset: + - class: SigAnaRecord + module_path: qlib.workflow.record_temp + kwargs: + ana_long_short: False + ann_scaler: 252 + - class: PortAnaRecord + module_path: qlib.workflow.record_temp + kwargs: + config: *port_analysis_config diff --git a/rdagent/scenarios/qlib/task_generator/model_template/read_exp_res.py b/rdagent/scenarios/qlib/task_generator/model_template/read_exp_res.py new file mode 100644 index 00000000..19930248 --- /dev/null +++ b/rdagent/scenarios/qlib/task_generator/model_template/read_exp_res.py @@ -0,0 +1,45 @@ +import pickle +from pathlib import Path + +import pandas as pd +import qlib +from mlflow.entities import ViewType +from mlflow.tracking import MlflowClient + +qlib.init() + +from qlib.workflow import R + +# here is the documents of the https://qlib.readthedocs.io/en/latest/component/recorder.html + +# TODO: list all the recorder and metrics + +# Assuming you have already listed the experiments +experiments = R.list_experiments() + +# Iterate through each experiment to find the latest recorder +experiment_name = None +latest_recorder = None +for experiment in experiments: + recorders = R.list_recorders(experiment_name=experiment) + for recorder_id in recorders: + if recorder_id is not None: + experiment_name = experiment + recorder = R.get_recorder(recorder_id=recorder_id, experiment_name=experiment) + end_time = recorder.info["end_time"] + if latest_recorder is None or end_time > latest_recorder.info["end_time"]: + latest_recorder = recorder + +# Check if the latest recorder is found +if latest_recorder is None: + print("No recorders found") +else: + print(f"Latest recorder: {latest_recorder}") + + # Load the specified file from the latest recorder + metrics = pd.Series(latest_recorder.list_metrics()) + + output_path = Path(__file__).resolve().parent / "qlib_res.csv" + metrics.to_csv(output_path) + + print(f"Output has been saved to {output_path}") diff --git a/rdagent/utils/env.py b/rdagent/utils/env.py index fb07f61e..02adf3fa 100644 --- a/rdagent/utils/env.py +++ b/rdagent/utils/env.py @@ -122,6 +122,7 @@ class DockerConf(BaseSettings): # And the extra data may be shared and the downloading can be time consuming. # So we just want to download it once. network: str | None = "bridge" # the network mode for the docker + shm_size: str | None = None class QlibDockerConf(DockerConf): @@ -131,6 +132,7 @@ class QlibDockerConf(DockerConf): mount_path: str = "/workspace/qlib_workspace/" default_entry: str = "qrun conf.yaml" extra_volumes: dict = {Path("~/.qlib/").expanduser().resolve(): "/root/.qlib/"} + shm_size: str | None = "16g" class DockerEnv(Env[DockerConf]): @@ -180,6 +182,7 @@ class DockerEnv(Env[DockerConf]): working_dir=self.conf.mount_path, # auto_remove=True, # remove too fast might cause the logs not to be get network=self.conf.network, + shm_size=self.conf.shm_size, ) logs = container.logs(stream=True) for log in logs: