diff --git a/fedot/api/api_utils/api_composer.py b/fedot/api/api_utils/api_composer.py index 934fe5cb46..0d295ce258 100644 --- a/fedot/api/api_utils/api_composer.py +++ b/fedot/api/api_utils/api_composer.py @@ -23,11 +23,13 @@ from fedot.core.pipelines.tuning.tuner_builder import TunerBuilder from fedot.core.repository.metrics_repository import MetricIDType from fedot.utilities.composer_timer import fedot_composer_timer +from fedot.core.context.context import ExecutionContext class ApiComposer: - def __init__(self, api_params: ApiParams, metrics: Union[MetricIDType, Sequence[MetricIDType]]): + def __init__(self, api_params: ApiParams, metrics: Union[MetricIDType, Sequence[MetricIDType]], + context: Optional[ExecutionContext] = None): self.log = default_log(self) self.params = api_params self.metrics = metrics @@ -39,6 +41,7 @@ def __init__(self, api_params: ApiParams, metrics: Union[MetricIDType, Sequence[ self.was_optimised = False # status flag indicating that tuner step was applied` self.was_tuned = False + self.context = context self.init_cache() def init_cache(self): @@ -155,6 +158,7 @@ def compose_pipeline(self, train_data: InputData, initial_assumption: Sequence[P .with_optimizer(self.params.get('optimizer')) .with_optimizer_params(parameters=self.params.optimizer_params) .with_metrics(self.metrics) + .with_context(self.context) .with_cache(self.operations_cache, self.preprocessing_cache, self.predictions_cache) .with_graph_generation_param(self.params.graph_generation_params) .build()) @@ -204,6 +208,7 @@ def tune_final_pipeline(self, train_data: InputData, .with_timeout(datetime.timedelta(minutes=tuner_plan.timeout_minutes)) .with_eval_time_constraint(self.params.composer_requirements.max_graph_fit_time) .with_requirements(self.params.composer_requirements) + .with_context(self.context) .build(train_data)) with self.timer.launch_tuning(): diff --git a/fedot/api/api_utils/api_params_repository.py b/fedot/api/api_utils/api_params_repository.py index c4d4430184..9dda08094a 100644 --- a/fedot/api/api_utils/api_params_repository.py +++ b/fedot/api/api_utils/api_params_repository.py @@ -25,9 +25,10 @@ class ApiParamsRepository: STATIC_INDIVIDUAL_METADATA_KEYS = {'use_input_preprocessing'} - def __init__(self, task_type: TaskTypesEnum): + def __init__(self, task_type: TaskTypesEnum, context: Optional[ExecutionContext] = None): self.task_type = task_type self.default_params = ApiParamsRepository.default_params_for_task(self.task_type) + self.context = context @staticmethod def default_params_for_task(task_type: TaskTypesEnum) -> dict: @@ -75,12 +76,21 @@ def get_params_for_gp_algorithm_params(self, params: dict) -> dict: if params.get('genetic_scheme') == 'steady_state': gp_algorithm_params['genetic_scheme_type'] = GeneticSchemeTypesEnum.steady_state - gp_algorithm_params['mutation_types'] = ApiParamsRepository._get_default_mutations(self.task_type, params) + gp_algorithm_params['mutation_types'] = ApiParamsRepository._get_default_mutations(self.task_type, params, + self.context) gp_algorithm_params['seed'] = params['seed'] return gp_algorithm_params + + @staticmethod + def _get_default_mutations(task_type: TaskTypesEnum, params, context: Optional[ExecutionContext] = None) -> Sequence[MutationTypesEnum]: + if context: + return context.default_mutations.get_default_mutation(task_type, params) + else: + return _get_default_mutations_core(task_type, params) + @staticmethod - def _get_default_mutations(task_type: TaskTypesEnum, params) -> Sequence[MutationTypesEnum]: + def _get_default_mutations_core(task_type: TaskTypesEnum, params) -> Sequence[MutationTypesEnum]: mutations = [parameter_change_mutation, MutationTypesEnum.single_change, MutationTypesEnum.single_drop, diff --git a/fedot/api/main.py b/fedot/api/main.py index e770a07b56..9a4770881a 100644 --- a/fedot/api/main.py +++ b/fedot/api/main.py @@ -51,6 +51,7 @@ from fedot.utilities.define_metric_by_task import MetricByTask from fedot.utilities.memory import MemoryAnalytics from fedot.utilities.project_import_export import export_project_to_zip, import_project_from_zip +from fedot.core.contex.context import ExecutionContext NOT_FITTED_ERR_MSG = 'Model not fitted yet' @@ -102,9 +103,15 @@ def __init__(self, logging_level: int = logging.ERROR, safe_mode: bool = False, n_jobs: int = -1, + context: Optional[str] = None, **composer_tuner_params ): + if context is None: + self.context = None + elif isinstance(context, str): + self.context = ExecutionContext(extension_name=context) + set_random_seed(seed) self.log = self._init_logger(logging_level) @@ -115,7 +122,7 @@ def __init__(self, passed_metrics = self.params.get('metric') self.metrics = ensure_wrapped_in_sequence(passed_metrics) if passed_metrics else default_metrics - self.api_composer = ApiComposer(self.params, self.metrics) + self.api_composer = ApiComposer(self.params, self.metrics, self.context) # Initialize data processors for data preprocessing and preliminary data analysis self.data_processor = ApiDataProcessor(task=self.params.task, @@ -299,7 +306,8 @@ def tune_tensordata(self, .with_n_jobs(common_tune_plan.n_jobs) .with_metric(common_tune_plan.metric) .with_iterations(iterations) - .with_timeout(timeout)) + .with_timeout(timeout) + .with_context(self.context)) pipeline_tuner = getattr(pipeline_tuner, tune_plan.builder_method_name)( tensor_data if tune_plan.use_tensor_runtime else common_tune_plan.input_data ) @@ -373,6 +381,7 @@ def tune(self, .with_metric(tune_plan.metric) .with_iterations(iterations) .with_timeout(timeout) + .with_context(self.context) .build(tune_input_data)) self.current_pipeline = pipeline_tuner.tune(self.current_pipeline, show_progress=show_progress) @@ -678,7 +687,8 @@ def get_metrics(self, data_producer=lambda: (yield self.train_data, self.test_data), validation_blocks=validation_blocks, eval_n_jobs=self.params.n_jobs, - do_unfit=False) + do_unfit=False, + context=self.context) metrics = obj_eval.evaluate(self.current_pipeline).values metrics = {metric_name: round(abs(metric), rounding_order) for (metric_name, metric) in diff --git a/fedot/core/composer/composer_builder.py b/fedot/core/composer/composer_builder.py index e4c37bcbb0..1081e41b7d 100644 --- a/fedot/core/composer/composer_builder.py +++ b/fedot/core/composer/composer_builder.py @@ -15,6 +15,7 @@ from fedot.core.caching.predictions_cache import PredictionsCache from fedot.core.composer.composer import Composer from fedot.core.composer.gp_composer.gp_composer import GPComposer +from fedot.core.context.context import ExecutionContext from fedot.core.optimisers.objective.metrics_objective import MetricsObjective from fedot.core.pipelines.pipeline import Pipeline from fedot.core.pipelines.pipeline_composer_requirements import PipelineComposerRequirements @@ -58,6 +59,13 @@ def __init__(self, task: Task): self.preprocessing_cache: Optional[PreprocessingCache] = None self.predictions_cache: Optional[PredictionsCache] = None + self.context: Optional[ExecutionContext] = None + + def with_context(self, context): + if context: + self.context = context + return self + def with_composer(self, composer_cls: Optional[Type[Composer]]): if composer_cls is not None: self.composer_cls = composer_cls diff --git a/fedot/core/composer/gp_composer/gp_composer.py b/fedot/core/composer/gp_composer/gp_composer.py index b4b7a7421a..705110480e 100644 --- a/fedot/core/composer/gp_composer/gp_composer.py +++ b/fedot/core/composer/gp_composer/gp_composer.py @@ -14,6 +14,7 @@ from fedot.core.composer.composer import Composer from fedot.core.data.data import InputData from fedot.core.data.multi_modal import MultiModalData +from fedot.core.context.context import ExecutionContext from fedot.core.optimisers.objective.data_objective_eval import ( PipelineObjectiveEvaluate, ) @@ -40,7 +41,8 @@ def __init__(self, optimizer: GraphOptimizer, composer_requirements: PipelineComposerRequirements, operations_cache: Optional[OperationsCache] = None, preprocessing_cache: Optional[PreprocessingCache] = None, - predictions_cache: Optional[PredictionsCache] = None): + predictions_cache: Optional[PredictionsCache] = None, + context: Optional[str] = None,): super().__init__(optimizer, composer_requirements) self.composer_requirements = composer_requirements self.operations_cache: Optional[OperationsCache] = operations_cache @@ -48,11 +50,14 @@ def __init__(self, optimizer: GraphOptimizer, self.predictions_cache: Optional[PredictionsCache] = predictions_cache self.best_models: Collection[Pipeline] = () + self.context = ExecutionContext(extension_name=context) def compose_pipeline(self, data: Union[InputData, MultiModalData]) -> Union[Pipeline, Sequence[Pipeline]]: # Define data source data_splitter = DataSourceSplitter(self.composer_requirements.cv_folds, shuffle=True) + if self.context: + data_splitter.build = self.context.data_source_splitter.build data_producer = data_splitter.build(data) parallelization_mode = self.composer_requirements.parallelization_mode @@ -71,7 +76,9 @@ def compose_pipeline(self, data: Union[InputData, MultiModalData]) -> Union[Pipe preprocessing_cache=self.preprocessing_cache, predictions_cache=self.predictions_cache, validation_blocks=data_splitter.validation_blocks, - eval_n_jobs=n_jobs_for_evaluation) + eval_n_jobs=n_jobs_for_evaluation, + context=context) + objective_function = objective_evaluator.evaluate # Define callback for computing intermediate metrics if needed diff --git a/fedot/core/context/__init__.py b/fedot/core/context/__init__.py new file mode 100644 index 0000000000..367791304e --- /dev/null +++ b/fedot/core/context/__init__.py @@ -0,0 +1 @@ +from .context import ExecutionContext \ No newline at end of file diff --git a/fedot/core/context/context.py b/fedot/core/context/context.py new file mode 100644 index 0000000000..b77dd3c995 --- /dev/null +++ b/fedot/core/context/context.py @@ -0,0 +1,114 @@ +from typing import Dict, Any, Optional, Callable +from fedot.extensions.registry import get_registered_extension, register_extension + + +class ExecutionContext: + def __init__(self, extension_name: str = "industrial", extra_params: Optional[Dict[str, Any]] = None): + self.extension_name = extension_name + self.extra_params = extra_params or {} + self._instances: Dict[str, Any] = {} + self._overridden: Dict[str, Any] = {} + + manifest = get_registered_extension(extension_name) + if manifest is None: + raise ValueError(f"Extension '{extension_name}' not registered") + self._manifest = manifest + + self._protocol_classes = self._manifest.protocols or {} + + def _get_protocol_class(self, protocol_name: str) -> Callable: + if protocol_name in self._protocol_classes: + return self._protocol_classes[protocol_name] + raise ValueError(f"No implementation for protocol '{protocol_name}'") + + def _get_instance(self, protocol_name: str) -> Any: + if protocol_name in self._overridden: + override = self._overridden[protocol_name] + if isinstance(override, type): + return override(**self.extra_params) + return override + + if protocol_name not in self._instances: + protocol_class = self._get_protocol_class(protocol_name) + self._instances[protocol_name] = protocol_class(**self.extra_params) + return self._instances[protocol_name] + + @property + def splitter(self): + return self._get_instance("splitter") + + @property + def data_merger(self): + return self._get_instance("data_merger") + + @property + def image_merger(self): + return self._get_instance("image_merger") + + @property + def ts_merger(self): + return self._get_instance("ts_merger") + + @property + def text_merger(self): + return self._get_instance("text_merger") + + @property + def tuner_class(self): + return self._get_instance("tuner_class") + + @property + def data_source_splitter(self): + return self._get_instance("data_source_splitter") + + @property + def operation_predict(self): + return self._get_instance("operation_predict") + + @property + def lagged_transformer(self): + return self._get_instance("lagged_transformer") + + @property + def topological_features(self): + return self._get_instance("topological_features") + + @property + def ts_smoothing(self): + return self._get_instance("ts_smoothing") + + @property + def api_composer_tune(self): + return self._get_instance("api_composer_tune") + + @property + def reproduction(self): + return self._get_instance("reproduction") + + @property + def search_space(self): + return self._get_instance("search_space") + + @property + def default_mutations(self): + return self._get_instance("default_mutations") + + @property + def evaluator(self): + return self._get_instance("evaluator") + + def __setattr__(self, name: str, value: Any) -> None: + if name in ('extra_params', '_instances', '_overridden', '_protocol_classes', + '_manifest', 'extension_name'): + super().__setattr__(name, value) + else: + self._overridden[name] = value + + def __getattr__(self, name: str): + if name in self._overridden: + return self._overridden[name] + + if name in ('_protocol_classes', '_instances', '_manifest'): + return super().__getattribute__(name) + + raise AttributeError(f"'{type(self).__name__}' object has no attribute '{name}'") diff --git a/fedot/core/context/industrial_backend.py b/fedot/core/context/industrial_backend.py new file mode 100644 index 0000000000..8ddbe9c01d --- /dev/null +++ b/fedot/core/context/industrial_backend.py @@ -0,0 +1,189 @@ +from typing import List, Optional, Union, Any, TYPE_CHECKING +import numpy as np + +if TYPE_CHECKING: + from fedot.core.data.data import InputData, OutputData + from fedot.core.data.multi_modal import MultiModalData + from fedot.core.pipelines.pipeline import Pipeline + from golem.core.optimisers.fitness import Fitness + from golem.core.tuning.optuna_tuner import OptunaTuner + from fedot.core.repository.tasks import TaskTypesEnum + from fedot.core.pipelines.tuning.tuner import BaseTuner + +from fedot.core.protocols.protocols import ( + SplitterProtocol, + DataMergerProtocol, + ImageMergerProtocol, + TSMergerProtocol, + TextMergerProtocol, + DataSourceSplitterProtocol, + TunerClassProtocol, + ReproductionProtocol, + EvaluatorProtocol, + SearchSpaceProtocol, + DefaultMutationsProtocol, + OperationPredictProtocol, + LaggedTransformerProtocol, + TopologicalFeaturesProtocol, + TsSmoothingProtocol, + ApiComposerTuneProtocol, +) + + +class IndustrialSplitter(SplitterProtocol): + def split_any(self, data: 'InputData', split_ratio: float, shuffle: bool, + stratify: bool, random_seed: int, **kwargs): + from fedot.industrial.core.repository.industrial_implementations.abstract import split_any_industrial + return split_any_industrial(data, split_ratio, shuffle, stratify, random_seed, **kwargs) + + def split_time_series(self, data: 'InputData', validation_blocks: Optional[int] = None, **kwargs): + from fedot.industrial.core.repository.industrial_implementations.abstract import split_time_series_industrial + return split_time_series_industrial(data, validation_blocks, **kwargs) + + +class IndustrialDataMerger(DataMergerProtocol): + @staticmethod + def get(outputs: List['OutputData']): + from fedot.industrial.core.repository.industrial_implementations.abstract import get_merger_industrial + return get_merger_industrial(outputs) + + def merge_predicts(self, predicts: List[np.ndarray]) -> np.ndarray: + from fedot.industrial.core.repository.industrial_implementations.abstract import merge_industrial_predicts + return merge_industrial_predicts(predicts) + + @staticmethod + def find_main_output(outputs: List['OutputData']) -> 'OutputData': + from fedot.industrial.core.repository.industrial_implementations.abstract import find_main_output_industrial + return find_main_output_industrial(outputs) + + +class IndustrialImageMerger(ImageMergerProtocol): + def preprocess_predicts(self, predicts: List[np.ndarray]) -> List[np.ndarray]: + from fedot.industrial.core.repository.industrial_implementations.abstract import preprocess_industrial_predicts + return preprocess_industrial_predicts(predicts) + + +class IndustrialTSMerger(TSMergerProtocol): + + def postprocess_predicts(self, merged: np.ndarray) -> np.ndarray: + from fedot.industrial.core.repository.industrial_implementations.abstract import postprocess_industrial_predicts + return postprocess_industrial_predicts(merged) + + +class IndustrialTextMerger(TextMergerProtocol): + def merge_predicts(self, predicts: List[np.ndarray]) -> np.ndarray: + from fedot.industrial.core.repository.industrial_implementations.abstract import merge_industrial_predicts + return merge_industrial_predicts(predicts) + + +class IndustrialDataSourceSplitterBuilder(DataSourceSplitterProtocol): + def build(self, data: Union['InputData', 'MultiModalData']): + from fedot.industrial.core.repository.industrial_implementations.abstract import build_industrial + return build_industrial(data) + + +class IndustrialTunerClass(TunerClassProtocol): + def __init__(self, **kwargs): + self.backend = kwargs.get("backend", "default") + + def optuna_tuner(self, objective_evaluate, task, iterations, max_lead_time=None, **kwargs): + from golem.core.tuning.optuna_tuner import OptunaTuner + from fedot.industrial.core.repository.industrial_implementations.ml_optimisation import DaskOptunaTuner + + if "dask" in self.backend: + return DaskOptunaTuner(objective_evaluate, task, iterations, max_lead_time, **kwargs) + else: + return OptunaTuner(objective_evaluate, task, iterations, max_lead_time, **kwargs) + + +class IndustrialReproduction(ReproductionProtocol): + def reproduce(self, population, evaluator, **kwargs): + from fedot.industrial.core.repository.industrial_implementations.optimisation import reproduce_industrial + return reproduce_industrial(population, evaluator, **kwargs) + + def reproduce_uncontrolled(self, population, **kwargs): + from fedot.industrial.core.repository.industrial_implementations.optimisation import \ + reproduce_controlled_industrial + return reproduce_controlled_industrial(population, **kwargs) + + +class IndustrialEvaluator(EvaluatorProtocol): + def evaluate(self, graph: 'Pipeline') -> 'Fitness': + from fedot.industrial.core.metrics.pipeline import industrial_evaluate_pipeline + return industrial_evaluate_pipeline(graph) + + +class IndustrialSearchSpace(SearchSpaceProtocol): + def get_parameters_dict(self): + from fedot.industrial.core.tuning.search_space import get_industrial_search_space + return get_industrial_search_space() + + +class IndustrialDefaultMutations(DefaultMutationsProtocol): + @staticmethod + def get_default_mutations(task_type: 'TaskTypesEnum', params): + from fedot.industrial.core.repository.industrial_implementations.optimisation import \ + _get_default_industrial_mutations + return _get_default_industrial_mutations(task_type, params) + + +class IndustrialOperationPredict(OperationPredictProtocol): + def predict(self, fitted_operation, data: 'InputData', params=None, output_mode='default'): + from fedot.industrial.core.repository.industrial_implementations.abstract import predict_industrial + return predict_industrial(fitted_operation, data, params, output_mode) + + def predict_for_fit(self, fitted_operation, data: 'InputData', params=None, output_mode='default'): + from fedot.industrial.core.repository.industrial_implementations.abstract import predict_for_fit_industrial + return predict_for_fit_industrial(fitted_operation, data, params, output_mode) + + def _predict(self, fitted_operation, data: 'InputData', params=None, output_mode='default', + is_fit_stage=False, predictions_cache=None, fold_id=None, descriptive_id=None): + from fedot.industrial.core.repository.industrial_implementations.abstract import predict_operation_industrial + return predict_operation_industrial(fitted_operation, data, params, output_mode, + is_fit_stage, predictions_cache, fold_id, descriptive_id) + + +class IndustrialLaggedTransformer(LaggedTransformerProtocol): + def _update_column_types(self, output_data: 'OutputData'): + from fedot.industrial.core.repository.industrial_implementations.data_transformation import \ + update_column_types_industrial + return update_column_types_industrial(output_data) + + def transform(self, input_data: 'InputData') -> 'OutputData': + from fedot.industrial.core.repository.industrial_implementations.data_transformation import \ + transform_lagged_industrial + return transform_lagged_industrial(input_data) + + def transform_for_fit(self, input_data: 'InputData') -> 'OutputData': + from fedot.industrial.core.repository.industrial_implementations.data_transformation import \ + transform_lagged_for_fit_industrial + return transform_lagged_for_fit_industrial(input_data) + + def _check_and_correct_window_size(self, time_series: np.ndarray, forecast_length: int): + from fedot.industrial.core.repository.industrial_implementations.data_transformation import \ + _check_and_correct_window_size_industrial + return _check_and_correct_window_size_industrial(time_series, forecast_length) + + +class IndustrialTopologicalFeatures(TopologicalFeaturesProtocol): + def fit(self, input_data: 'InputData'): + from fedot.industrial.core.repository.industrial_implementations.abstract import fit_topo_extractor_industrial + return fit_topo_extractor_industrial(input_data) + + def transform(self, input_data: 'InputData') -> np.ndarray: + from fedot.industrial.core.repository.industrial_implementations.abstract import \ + transform_topo_extractor_industrial + return transform_topo_extractor_industrial(input_data) + + +class IndustrialTsSmoothing(TsSmoothingProtocol): + def transform(self, input_data: 'InputData') -> 'OutputData': + from fedot.industrial.core.repository.industrial_implementations.data_transformation import \ + transform_smoothing_industrial + return transform_smoothing_industrial(input_data) + + +class IndustrialApiComposerTune(ApiComposerTuneProtocol): + def tune_pipeline(self, train_data: 'InputData', pipeline: 'Pipeline', execution_plan=None) -> 'Pipeline': + from fedot.industrial.core.repository.industrial_implementations.ml_optimisation import tune_pipeline_industrial + return tune_pipeline_industrial(train_data, pipeline, execution_plan) diff --git a/fedot/core/context/industrial_manifest.py b/fedot/core/context/industrial_manifest.py new file mode 100644 index 0000000000..36a1439450 --- /dev/null +++ b/fedot/core/context/industrial_manifest.py @@ -0,0 +1,47 @@ +from fedot.extensions.contracts import ExtensionManifest +from fedot.extensions.registry import get_registered_extension, register_extension +from fedot.core.context.industrial_backend import ( + IndustrialSplitter, + IndustrialDataMerger, + IndustrialImageMerger, + IndustrialTSMerger, + IndustrialTextMerger, + IndustrialDataSourceSplitterBuilder, + IndustrialTunerClass, + IndustrialReproduction, + IndustrialEvaluator, + IndustrialSearchSpace, + IndustrialDefaultMutations, + IndustrialOperationPredict, + IndustrialLaggedTransformer, + IndustrialTopologicalFeatures, + IndustrialTsSmoothing, + IndustrialApiComposerTune, +) + +FEDOT_INDUSTRIAL_MANIFEST = ExtensionManifest( + name="industrial", + version="1.0.0", + models=(), + description="Industrial extension for FEDOT.", + protocols={ + "splitter": IndustrialSplitter, + "data_merger": IndustrialDataMerger, + "image_merger": IndustrialImageMerger, + "ts_merger": IndustrialTSMerger, + "text_merger": IndustrialTextMerger, + "data_source_splitter": IndustrialDataSourceSplitterBuilder, + "tuner_class": IndustrialTunerClass, + "reproduction": IndustrialReproduction, + "evaluator": IndustrialEvaluator, + "search_space": IndustrialSearchSpace, + "default_mutations": IndustrialDefaultMutations, + "operation_predict": IndustrialOperationPredict, + "lagged_transformer": IndustrialLaggedTransformer, + "topological_features": IndustrialTopologicalFeatures, + "ts_smoothing": IndustrialTsSmoothing, + "api_composer_tune": IndustrialApiComposerTune, + } +) + +register_extension(FEDOT_INDUSTRIAL_MANIFEST) diff --git a/fedot/core/data/data_split.py b/fedot/core/data/data_split.py index a000c6e46b..a18221dfa0 100644 --- a/fedot/core/data/data_split.py +++ b/fedot/core/data/data_split.py @@ -8,6 +8,7 @@ from fedot.core.data.multi_modal import MultiModalData from fedot.core.repository.dataset_types import DataTypesEnum from fedot.core.repository.tasks import TaskTypesEnum +from fedot.core.context.context import ExecutionContext def _split_input_data_by_indexes(origin_input_data: Union[InputData, MultiModalData], @@ -173,7 +174,8 @@ def train_test_data_setup(data: Union[InputData, MultiModalData], shuffle_flag: bool = False, stratify: bool = True, random_seed: int = 42, - validation_blocks: Optional[int] = None) -> Tuple[Union[InputData, MultiModalData], + validation_blocks: Optional[int] = None, + context: Optional[ExecutionContext] = None) -> Tuple[Union[InputData, MultiModalData], Union[InputData, MultiModalData]]: """ Function for train and test split for both InputData and MultiModalData @@ -203,11 +205,18 @@ def train_test_data_setup(data: Union[InputData, MultiModalData], 'random_seed': random_seed, 'validation_blocks': validation_blocks} if isinstance(data, InputData): - split_func_dict = {DataTypesEnum.multi_ts: _split_time_series, - DataTypesEnum.ts: _split_time_series, - DataTypesEnum.table: _split_any, - DataTypesEnum.image: _split_any, - DataTypesEnum.text: _split_any} + if context: + split_func_dict = {DataTypesEnum.multi_ts: context.splitter.split_time_series, + DataTypesEnum.ts: context.splitter.split_time_series, + DataTypesEnum.table: context.splitter.split_any, + DataTypesEnum.image: context.splitter.split_any, + DataTypesEnum.text: context.splitter.split_any,} + else: + split_func_dict = {DataTypesEnum.multi_ts: _split_time_series, + DataTypesEnum.ts: _split_time_series, + DataTypesEnum.table: _split_any, + DataTypesEnum.image: _split_any, + DataTypesEnum.text: _split_any} if data.data_type not in split_func_dict: raise TypeError((f'Unknown data type {type(data)}. Supported data types:' diff --git a/fedot/core/data/merge/data_merger.py b/fedot/core/data/merge/data_merger.py index a1dc312f0b..c4cfb1b150 100644 --- a/fedot/core/data/merge/data_merger.py +++ b/fedot/core/data/merge/data_merger.py @@ -9,6 +9,7 @@ from fedot.core.data.data import OutputData, InputData from fedot.core.data.merge.supplementary_data_merger import SupplementaryDataMerger from fedot.core.repository.dataset_types import DataTypesEnum +from fedot.core.context.context import ExecutionContext class DataMerger: @@ -23,10 +24,11 @@ class DataMerger: :param outputs: list with OutputData from parent nodes for merging """ - def __init__(self, outputs: List['OutputData'], data_type: DataTypesEnum = None): + def __init__(self, outputs: List['OutputData'], data_type: DataTypesEnum = None, context: Optional[ExecutionContext] = None,): self.log = default_log(self) self.outputs = outputs self.data_type = data_type or DataMerger.get_datatype_for_merge(output.data_type for output in outputs) + self.context = context # Ensure outputs are of equal length, find common index if it is not idx_list = [np.asarray(output.idx) for output in outputs] @@ -35,10 +37,10 @@ def __init__(self, outputs: List['OutputData'], data_type: DataTypesEnum = None) raise ValueError('There are no common indices for outputs') # Find first output with the main target & resulting task - self.main_output = DataMerger.find_main_output(outputs) + self.main_output = DataMerger.find_main_output(outputs, self.context) @staticmethod - def get(outputs: List['OutputData']) -> 'DataMerger': + def get_core(outputs: List['OutputData']) -> 'DataMerger': """ Construct appropriate data merger for the outputs. """ # Ensure outputs can be merged @@ -116,7 +118,7 @@ def preprocess_predicts(self, predicts: List[np.array]) -> List[np.array]: """ Pre-process (e.g. equalizes sizes, reshapes) and return list of arrays that can be merged. """ return list(map(atleast_2d, predicts)) - def merge_predicts(self, predicts: List[np.array]) -> np.array: + def merge_predicts_core(self, predicts: List[np.array]) -> np.array: # Finally, merge predictions into features for the next stage return np.concatenate(predicts, axis=-1) @@ -137,7 +139,7 @@ def is_forecast_index(output: 'OutputData'): return len(output.idx) != len(output.predict) @staticmethod - def find_main_output(outputs: List['OutputData']) -> 'OutputData': + def find_main_output_core(outputs: List['OutputData']) -> 'OutputData': """ Returns first output with main target or (if there are no main targets) the output with priority secondary target. """ priority_output = next((output for output in outputs @@ -147,11 +149,35 @@ def find_main_output(outputs: List['OutputData']) -> 'OutputData': i_priority_secondary = np.argmin(flow_lengths) priority_output = outputs[i_priority_secondary] return priority_output + @staticmethod + def find_main_output(outputs: List['OutputData'], context: Optional[ExecutionContext] = None) -> 'OutputData': + if context: + return context.data_merger.find_main_output(outputs) + else: + return find_main_output_core(outputs) + @staticmethod + def get(outputs: List['OutputData'], context: Optional[ExecutionContext] = None) -> 'DataMerger': + if context: + return context.data_merger.get(outputs) + else: + return get_core(outputs) + + def merge_predicts(self, predicts: List[np.array]) -> np.array: + if self.context: + return self.context.data_merger.merge_predicts(predicts) + else: + return self.merge_predicts_core(predicts) class ImageDataMerger(DataMerger): def preprocess_predicts(self, predicts: List[np.array]) -> List[np.array]: + if self.context: + return self.context.image_merger.preprocess_predicts(predicts) + else: + return self.preprocess_predicts_core(predicts) + + def preprocess_predicts_core(self, predicts: List[np.array]) -> List[np.array]: # Reshape predicts to 4d (idx, width, height, channels) reshaped_predicts = list(map(atleast_4d, predicts)) @@ -167,6 +193,12 @@ def preprocess_predicts(self, predicts: List[np.array]) -> List[np.array]: class TSDataMerger(DataMerger): def postprocess_predicts(self, merged_predicts: np.array) -> np.array: + if self.context: + return self.context.ts_merger.postprocess_predicts(merged_predicts) + else: + return self.postprocess_predicts_core(merged_predicts) + + def postprocess_predicts_core(self, merged_predicts: np.array) -> np.array: # Ensure that 1d-column timeseries remains 1d timeseries return flatten_extra_dim(merged_predicts) diff --git a/fedot/core/operations/evaluation/operation_implementations/data_operations/topological/fast_topological_extractor.py b/fedot/core/operations/evaluation/operation_implementations/data_operations/topological/fast_topological_extractor.py index cf734fa8e0..584aaced1b 100644 --- a/fedot/core/operations/evaluation/operation_implementations/data_operations/topological/fast_topological_extractor.py +++ b/fedot/core/operations/evaluation/operation_implementations/data_operations/topological/fast_topological_extractor.py @@ -3,6 +3,7 @@ from typing import Optional import numpy as np +from fedot.core.context.context import ExecutionContext try: from gph import ripser_parallel as ripser @@ -20,7 +21,7 @@ class TopologicalFeaturesImplementation(DataOperationImplementation): - def __init__(self, params: Optional[OperationParameters] = None): + def __init__(self, params: Optional[OperationParameters] = None, context: Optional[ExecutionContext] = None): super().__init__(params) self.window_size_as_share = params.get('window_size_as_share') self.max_homology_dimension = params.get('max_homology_dimension') @@ -30,15 +31,23 @@ def __init__(self, params: Optional[OperationParameters] = None): self.quantiles = (0.1, 0.25, 0.5, 0.75, 0.9) self._shape = len(self.quantiles) self._window_size = None + self.context = context def fit(self, input_data: InputData): + if self.context: + return self.context.topological_feature.fit(input_data) + else: + return self.fit_core(input_data) + + + def fit_core(self, input_data: InputData): self._window_size = int(input_data.features.shape[1] * self.window_size_as_share) self._window_size = max(self._window_size, 2) self._window_size = min(self._window_size, input_data.features.shape[1] - 2) self._window_size = max(self._window_size, 1) return self - def transform(self, input_data: InputData) -> OutputData: + def transform_core(self, input_data: InputData) -> OutputData: features = input_data.features with Parallel(n_jobs=self.n_jobs, prefer='processes') as parallel: topological_features = parallel(delayed(self._extract_features) @@ -52,6 +61,12 @@ def transform(self, input_data: InputData) -> OutputData: np.nan_to_num(result, copy=False, nan=0, posinf=0, neginf=0) return result + def transform(self, input_data: InputData): + if self.context: + return self.context.topological_feature.transform(input_data) + else: + return self.transform_core(input_data) + def _extract_features(self, x): x_sliced = np.array([x[i:self._window_size + i] for i in range(x.shape[0] - self._window_size + 1)]) x_processed = ripser(x_sliced, diff --git a/fedot/core/operations/evaluation/operation_implementations/data_operations/ts_transformations.py b/fedot/core/operations/evaluation/operation_implementations/data_operations/ts_transformations.py index 7222872be4..ea835d43fe 100644 --- a/fedot/core/operations/evaluation/operation_implementations/data_operations/ts_transformations.py +++ b/fedot/core/operations/evaluation/operation_implementations/data_operations/ts_transformations.py @@ -17,10 +17,10 @@ from fedot.core.operations.operation_parameters import OperationParameters from fedot.core.repository.dataset_types import DataTypesEnum from fedot.preprocessing.data_types import TYPE_TO_ID - +from fedot.core.context.context import ExecutionContext class LaggedImplementation(DataOperationImplementation): - def __init__(self, params: Optional[OperationParameters]): + def __init__(self, params: Optional[OperationParameters], context: Optional[ExecutionContext] = None): super().__init__(params) self.window_size_minimum = None @@ -30,6 +30,7 @@ def __init__(self, params: Optional[OperationParameters]): # Define logger object self.log = default_log(self) + self.context = context @property def window_size(self) -> Optional[int]: @@ -51,7 +52,13 @@ def fit(self, input_data): pass - def transform(self, input_data: InputData) -> OutputData: + def tranfsorm(self, input_data: InputData) -> OutputData: + if self.context: + return self.context.lagged_transformer.transform(input_data) + else: + return self.transform_core(input_data) + + def transform_core(self, input_data: InputData) -> OutputData: """ Method for transformation of time series to lagged form for predict stage Args: @@ -75,7 +82,7 @@ def transform(self, input_data: InputData) -> OutputData: self._update_column_types(output_data) return output_data - def transform_for_fit(self, input_data: InputData) -> OutputData: + def transform_for_fit_core(self, input_data: InputData) -> OutputData: """Method for transformation of time series to lagged form for fit stage Args: @@ -106,7 +113,14 @@ def transform_for_fit(self, input_data: InputData) -> OutputData: self._update_column_types(output_data) return output_data - def _check_and_correct_window_size(self, time_series: np.ndarray, forecast_length: int): + def transform_for_fit(self, input_data: InputData) -> OutputData: + if self.context: + return self.context.lagged_transformer.transform_for_fit(input_data) + else: + return self.transform_for_fit_core(input_data) + + + def _check_and_correct_window_size_core(self, time_series: np.ndarray, forecast_length: int): """ Method check if the length of the time series is not enough for lagged transformation @@ -142,7 +156,19 @@ def _check_and_correct_window_size(self, time_series: np.ndarray, forecast_lengt f"from {self.params.get('window_size')} to {self.window_size_minimum}")) self.params.update(window_size=self.window_size_minimum) + def _check_and_correct_window_size(self, time_series: np.ndarray, forecast_length: int): + if self.context: + return self.context.lagged_transformer._check_and_correct_window_size(time_series, forecast_length) + else: + return self._check_and_correct_window_size_core(time_series, forecast_length) + def _update_column_types(self, output_data: OutputData): + if self.context: + return self.context.lagged_transformer._update_column_types(output_data) + else: + return self._update_column_types_core(output_data) + + def _update_column_types_core(self, output_data: OutputData): """Update column types after lagged transformation. All features becomes ``float`` """ @@ -371,8 +397,9 @@ def __init__(self, params: Optional[OperationParameters]): class TsSmoothingImplementation(DataOperationImplementation): - def __init__(self, params: Optional[OperationParameters]): + def __init__(self, params: Optional[OperationParameters], context: Optional[ExecutionContext] = None): super().__init__(params) + self.context = context @property def window_size(self) -> int: @@ -387,7 +414,13 @@ def fit(self, input_data: InputData): pass - def transform(self, input_data: InputData) -> OutputData: + def transform(self, input_data: InputData): + if self.context: + self.context.ts_smoothing.transform(input_data) + else: + return self.transform_core(input_data) + + def transform_core(self, input_data: InputData) -> OutputData: """Method for smoothing time series Args: diff --git a/fedot/core/operations/operation.py b/fedot/core/operations/operation.py index faf2fdf341..18e5ecd8ed 100644 --- a/fedot/core/operations/operation.py +++ b/fedot/core/operations/operation.py @@ -27,7 +27,7 @@ class Operation: operation_type: name of the operation """ - def __init__(self, operation_type: str, **kwargs): + def __init__(self, operation_type: str, context: Optional[ExecutionContext] = None, **kwargs): self.operation_type = operation_type self._eval_strategy = None @@ -35,6 +35,7 @@ def __init__(self, operation_type: str, **kwargs): self.fitted_operation = None self.log = default_log(self) + self.context = context def _init(self, task: Task, **kwargs): params = kwargs.get('params') @@ -110,6 +111,48 @@ def predict(self, predictions_cache: Optional[PredictionsCache] = None, fold_id: Optional[int] = None, descriptive_id: Optional[str] = None): + if self.context: + return self.context.operation_predict.predict(fitted_operation, data, params, output_mode, predictions_cache, fold_id, descriptive_id) + else: + return self.predict_core(fitted_operation, data, params, output_mode, predictions_cache, fold_id, descriptive_id) + + def predict_for_fit(self, + fitted_operation, + data: InputData, + params: Optional[OperationParameters] = None, + output_mode: str = 'default', + predictions_cache: Optional[PredictionsCache] = None, + fold_id: Optional[int] = None, + descriptive_id: Optional[str] = None): + if self.context: + return self.context.operation_predict.predict_for_fit(fitted_operation, data, params, output_mode, predictions_cache, fold_id, descriptive_id) + else: + return self.predict_for_fit_core(fitted_operation, data, params, output_mode, predictions_cache, fold_id, descriptive_id) + + def _predict(self, + fitted_operation, + data: InputData, + params: Optional[OperationParameters] = None, + output_mode: str = 'default', + is_fit_stage: bool = False, + predictions_cache: Optional[PredictionsCache] = None, + fold_id: Optional[int] = None, + descriptive_id: Optional[str] = None): + if self.context: + return self.context.operation_predict._predict(fitted_operation, data, params, output_mode, + is_fit_stage, predictions_cache, fold_id, descriptive_id) + else: + return self._predict_core(fitted_operation, data, params, output_mode, is_fit_stage, + predictions_cache, fold_id, descriptive_id) + + def predict_core(self, + fitted_operation, + data: InputData, + params: Optional[Union[OperationParameters, dict]] = None, + output_mode: str = 'default', + predictions_cache: Optional[PredictionsCache] = None, + fold_id: Optional[int] = None, + descriptive_id: Optional[str] = None): """This method is used for defining and running of the evaluation strategy to predict with the data provided @@ -120,10 +163,10 @@ def predict(self, output_mode: string with information about output of operation, for example, is the operation predict probabilities or class labels """ - return self._predict(fitted_operation, data, params, output_mode, is_fit_stage=False, + return self._predict_core(fitted_operation, data, params, output_mode, is_fit_stage=False, predictions_cache=predictions_cache, fold_id=fold_id, descriptive_id=descriptive_id) - def predict_for_fit(self, + def predict_for_fit_core(self, fitted_operation, data: InputData, params: Optional[OperationParameters] = None, @@ -144,7 +187,7 @@ def predict_for_fit(self, return self._predict(fitted_operation, data, params, output_mode, is_fit_stage=True, predictions_cache=predictions_cache, fold_id=fold_id, descriptive_id=descriptive_id) - def _predict(self, + def _predict_core(self, fitted_operation, data: InputData, params: Optional[OperationParameters] = None, diff --git a/fedot/core/optimisers/objective/data_objective_eval.py b/fedot/core/optimisers/objective/data_objective_eval.py index f139bee5c4..f547127c40 100644 --- a/fedot/core/optimisers/objective/data_objective_eval.py +++ b/fedot/core/optimisers/objective/data_objective_eval.py @@ -15,6 +15,7 @@ from fedot.core.operations.model import Model from fedot.core.pipelines.pipeline import Pipeline from fedot.utilities.debug import is_recording_mode, save_debug_info_for_pipeline +from fedot.core.context.context import ExecutionContext DataSource = Callable[[], Iterable[Tuple[InputData, InputData]]] @@ -44,7 +45,9 @@ def __init__(self, preprocessing_cache: Optional[PreprocessingCache] = None, predictions_cache: Optional[PredictionsCache] = None, eval_n_jobs: int = 1, - do_unfit: bool = True): + do_unfit: bool = True, + context: Optional[ExecutionContext] = None + ): super().__init__(objective, eval_n_jobs=eval_n_jobs) self._data_producer = data_producer self._time_constraint = time_constraint @@ -52,12 +55,21 @@ def __init__(self, self._operations_cache = operations_cache self._preprocessing_cache = preprocessing_cache self._predictions_cache = predictions_cache + self.context = context self._log = default_log(self) self._do_unfit = do_unfit def evaluate(self, graph: Pipeline) -> Fitness: # Seems like a workaround for situation when logger is lost # when adapting and restoring it to/from OptGraph. + + if self.context: + return self.context.evaluator.evaluate(graph) + + else: + return self.evaluation_core(graph) + + def evaluation_core(self, graph: Pipeline) -> Fitness: graph.log = self._log graph_id = graph.root_node.descriptive_id diff --git a/fedot/core/optimisers/objective/data_source_splitter.py b/fedot/core/optimisers/objective/data_source_splitter.py index c004b86b23..0a31428f6f 100644 --- a/fedot/core/optimisers/objective/data_source_splitter.py +++ b/fedot/core/optimisers/objective/data_source_splitter.py @@ -39,7 +39,8 @@ def __init__(self, split_ratio: Optional[float] = None, shuffle: bool = False, stratify: bool = True, - random_seed: int = 42): + random_seed: int = 42, + context: Optional[ExecutionContext] = None): self.cv_folds = cv_folds self.validation_blocks = validation_blocks self.split_ratio = split_ratio @@ -47,12 +48,19 @@ def __init__(self, self.stratify = stratify self.random_seed = random_seed self.log = default_log(self) + self.context = context def build_tensordata(self, tensor_data) -> DataSource: input_data = tensordata_to_input_data(tensor_data) return self.build(input_data) def build(self, data: Union[InputData, MultiModalData]) -> DataSource: + if self.context: + return self.context.data_source_splitter.build(data) + else: + return self.build_core(data) + + def build_core(self, data: Union[InputData, MultiModalData]) -> DataSource: # define split_ratio self.split_ratio = self.split_ratio or default_data_split_ratio_by_task[data.task.task_type] diff --git a/fedot/core/pipelines/tuning/search_space.py b/fedot/core/pipelines/tuning/search_space.py index 2724175253..70b594ce29 100644 --- a/fedot/core/pipelines/tuning/search_space.py +++ b/fedot/core/pipelines/tuning/search_space.py @@ -6,6 +6,8 @@ from fedot.core.utils import NESTED_PARAMS_LABEL +from fedot.core.context.context import ExecutionContext + class PipelineSearchSpace(SearchSpace): """ @@ -18,13 +20,21 @@ class PipelineSearchSpace(SearchSpace): def __init__(self, custom_search_space: Optional[OperationParametersMapping] = None, - replace_default_search_space: bool = False): + replace_default_search_space: bool = False, + context: Optional[ExecutionContext] = None): self.custom_search_space = custom_search_space self.replace_default_search_space = replace_default_search_space + self.context = context parameters_per_operation = self.get_parameters_dict() super().__init__(parameters_per_operation) def get_parameters_dict(self): + if self.context: + return self.context.search_space.get_parameters_dict() + else: + return self.get_parameters_dict_core() + + def get_parameters_dict_core(self): parameters_per_operation = { 'kmeans': { 'n_clusters': { diff --git a/fedot/core/pipelines/tuning/tuner_builder.py b/fedot/core/pipelines/tuning/tuner_builder.py index 3d9e33b7e9..6ae9c2c953 100644 --- a/fedot/core/pipelines/tuning/tuner_builder.py +++ b/fedot/core/pipelines/tuning/tuner_builder.py @@ -37,6 +37,11 @@ def __init__(self, task: Task): self.eval_time_constraint = None self.additional_params = {} self.adapter = PipelineAdapter() + self.context: Optional[ExecutionContext] = None + + def with_context(self, context: ExecutionContext): # ← добавить метод + self.context = context + return self def with_tuner(self, tuner: Type[BaseTuner]): self.tuner_class = tuner @@ -110,6 +115,7 @@ def _build_tuner(self, data_producer, validation_blocks: int) -> BaseTuner: time_constraint=self.eval_time_constraint, eval_n_jobs=self.n_jobs, # because tuners are not parallelized validation_blocks=validation_blocks, + context=context ) tuner = self.tuner_class(objective_evaluate=objective_evaluate, adapter=self.adapter, @@ -122,11 +128,11 @@ def _build_tuner(self, data_producer, validation_blocks: int) -> BaseTuner: return tuner def build(self, data: InputData) -> BaseTuner: - data_splitter = DataSourceSplitter(self.cv_folds, validation_blocks=self.validation_blocks) + data_splitter = DataSourceSplitter(self.cv_folds, validation_blocks=self.validation_blocks, context=self.context) data_producer = data_splitter.build(data) return self._build_tuner(data_producer, data_splitter.validation_blocks) def build_tensordata(self, tensor_data) -> BaseTuner: - data_splitter = DataSourceSplitter(self.cv_folds, validation_blocks=self.validation_blocks) + data_splitter = DataSourceSplitter(self.cv_folds, validation_blocks=self.validation_blocks, context=self.context) data_producer = data_splitter.build_tensordata(tensor_data) return self._build_tuner(data_producer, data_splitter.validation_blocks) diff --git a/fedot/core/protocols/__init__.py b/fedot/core/protocols/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/fedot/core/protocols/protocols.py b/fedot/core/protocols/protocols.py new file mode 100644 index 0000000000..726d16ccd0 --- /dev/null +++ b/fedot/core/protocols/protocols.py @@ -0,0 +1,160 @@ +from typing import Any, List, Protocol, Optional, Callable, Union, Tuple, Sequence, TYPE_CHECKING +import numpy as np + +from fedot.core.data.data import InputData, OutputData +from fedot.core.repository.tasks import TaskTypesEnum +from golem.core.optimisers.fitness import Fitness +from golem.core.tuning.tuner_interface import BaseTuner + +if TYPE_CHECKING: + from fedot.core.pipelines.pipeline import Pipeline + from fedot.core.data.multi_modal import MultiModalData + from fedot.core.data.merge.data_merger import DataMerger + +class EvaluatorProtocol(Protocol): + def evaluate(self, graph: 'Pipeline') -> Fitness: + ... + + +class SearchSpaceProtocol(Protocol): + def get_parameters_dict(self) -> dict: + ... + + +class DefaultMutationsProtocol(Protocol): + @staticmethod + def __call__(task_type: TaskTypesEnum, params: Any) -> Sequence[Any]: + ... + + +class DataMergerProtocol(Protocol): + @staticmethod + def get(outputs: List[OutputData]) -> 'DataMerger': + ... + + def merge_predicts(self, predicts: List[np.ndarray]) -> np.ndarray: + ... + + @staticmethod + def find_main_output(outputs: List[OutputData]) -> OutputData: + ... + + def preprocess_predicts(self, predicts: List[np.ndarray]) -> List[np.ndarray]: + ... + + def postprocess_predicts(self, merged: np.ndarray) -> np.ndarray: + ... + + +class ImageMergerProtocol(Protocol): + def preprocess_predicts(self, predicts: List[np.ndarray]) -> List[np.ndarray]: + ... + + def merge_predicts(self, predicts: List[np.ndarray]) -> np.ndarray: + ... + + +class TSMergerProtocol(Protocol): + def merge_predicts(self, predicts: List[np.ndarray]) -> np.ndarray: + ... + + def merge_targets(self, targets: List[np.ndarray]) -> np.ndarray: + ... + + def preprocess_predicts(self, predicts: List[np.ndarray]) -> List[np.ndarray]: + ... + + def postprocess_predicts(self, merged: np.ndarray) -> np.ndarray: + ... + + +class TextMergerProtocol(Protocol): + def merge_predicts(self, predicts: List[np.ndarray]) -> np.ndarray: + ... + + def postprocess_predicts(self, merged: np.ndarray) -> np.ndarray: + ... + + +class DataSourceSplitterProtocol(Protocol): + def build(self, data: Union[InputData, 'MultiModalData']) -> Callable: + ... + + +class SplitterProtocol(Protocol): + def split_any(self, data: InputData, split_ratio: float, shuffle: bool, + stratify: bool, random_seed: int, **kwargs) -> Tuple[InputData, InputData]: + ... + + def split_time_series(self, data: InputData, validation_blocks: Optional[int] = None, + **kwargs) -> Tuple[InputData, InputData]: + ... + + +class OperationPredictProtocol(Protocol): + def predict(self, fitted_operation, data: InputData, + params: Optional[Any] = None, output_mode: str = 'default') -> OutputData: + ... + + def predict_for_fit(self, fitted_operation, data: InputData, + params: Optional[Any] = None, output_mode: str = 'default') -> OutputData: + ... + + def _predict(self, fitted_operation, data: InputData, params: Optional[Any] = None, + output_mode: str = 'default', is_fit_stage: bool = False, + predictions_cache: Optional[Any] = None, fold_id: Optional[int] = None, + descriptive_id: Optional[str] = None) -> OutputData: + ... + + +class LaggedTransformerProtocol(Protocol): + def _update_column_types(self, output_data: OutputData) -> None: + ... + + def transform(self, input_data: InputData) -> OutputData: + ... + + def transform_for_fit(self, input_data: InputData) -> OutputData: + ... + + def _check_and_correct_window_size(self, time_series: np.ndarray, forecast_length: int) -> None: + ... + + +class TopologicalFeaturesProtocol(Protocol): + def fit(self, input_data: InputData) -> Any: + ... + + def transform(self, input_data: InputData) -> np.ndarray: + ... + + +class TsSmoothingProtocol(Protocol): + def transform(self, input_data: InputData) -> OutputData: + ... + + +class TunerClassProtocol(Protocol): + def __call__(self, objective_evaluate: Any, task: Any, iterations: int, + max_lead_time: Optional[float] = None, **kwargs) -> BaseTuner: + ... + + +class ApiComposerTuneProtocol(Protocol): + def __call__(self, train_data: InputData, pipeline: 'Pipeline', + execution_plan: Optional[Any] = None) -> 'Pipeline': + ... + + +class ReproductionProtocol(Protocol): + def reproduce(self, population: List[Any], evaluator: Any, **kwargs) -> List[Any]: + ... + + def reproduce_uncontrolled(self, population: List[Any], **kwargs) -> List[Any]: + ... + + +class VerificationRulesProtocol(Protocol): + class_rules: List[Callable] + ts_rules: List[Callable] + common_rules: List[Callable] \ No newline at end of file diff --git a/fedot/extensions/contracts.py b/fedot/extensions/contracts.py index 3b3917866d..ab10c38418 100644 --- a/fedot/extensions/contracts.py +++ b/fedot/extensions/contracts.py @@ -37,6 +37,7 @@ class ExtensionManifest: version: str models: Tuple[ExternalModelSpec, ...] module: Optional[str] = None + protocols: Optional[Dict[str, Callable[..., Any]]] = None description: str = '' diff --git a/fedot/extensions/registry.py b/fedot/extensions/registry.py index 43920c5045..63c19f858d 100644 --- a/fedot/extensions/registry.py +++ b/fedot/extensions/registry.py @@ -163,3 +163,4 @@ def smoke_test_extension(manifest: ExtensionManifest): message=f'Factory for model "{model.name}" returned None.')) return Right(manifest) + diff --git a/fedot/industrial/core/repository/initializer_industrial_models.py b/fedot/industrial/core/repository/initializer_industrial_models.py index eeb00a35d3..2f2c737b12 100644 --- a/fedot/industrial/core/repository/initializer_industrial_models.py +++ b/fedot/industrial/core/repository/initializer_industrial_models.py @@ -39,70 +39,8 @@ from fedot.industrial.core.repository.model_repository import overload_model_implementation from fedot.industrial.core.tuning.search_space import get_industrial_search_space -FEDOT_METHOD_TO_REPLACE = [(PipelineObjectiveEvaluate, "evaluate"), - (PipelineSearchSpace, "get_parameters_dict"), - (ApiParamsRepository, "_get_default_mutations"), - (DataMerger, "find_main_output"), - (DataMerger, "get"), - (DataMerger, "merge_predicts"), - (ImageDataMerger, "preprocess_predicts"), - (ImageDataMerger, "merge_predicts"), - (TSDataMerger, "merge_predicts"), - (TSDataMerger, "merge_targets"), - (TSDataMerger, 'postprocess_predicts'), - (TSDataMerger, 'preprocess_predicts'), - (DataSourceSplitter, "build"), - (fedot_data_split, "_split_any"), - (fedot_data_split, "_split_time_series"), - (Operation, "_predict"), - (Operation, "predict"), - (Operation, "predict_for_fit"), - (LaggedImplementation, '_update_column_types'), - (LaggedImplementation, 'transform'), - (TopologicalFeaturesImplementation, 'fit'), - (TopologicalFeaturesImplementation, 'transform'), - (LaggedImplementation, 'transform_for_fit'), - (LaggedImplementation, '_check_and_correct_window_size'), - (TsSmoothingImplementation, 'transform'), - (OptunaImpl, 'OptunaTuner'), - (ApiComposer, 'tune_final_pipeline'), - (ReproductionController, 'reproduce_uncontrolled'), - (ReproductionController, 'reproduce')] -INDUSTRIAL_REPLACE_METHODS = [industrial_evaluate_pipeline, - get_industrial_search_space, - _get_default_industrial_mutations, - find_main_output_industrial, - get_merger_industrial, - merge_industrial_predicts, - preprocess_industrial_predicts, - merge_industrial_predicts, - merge_industrial_predicts, - merge_industrial_targets, - postprocess_industrial_predicts, - preprocess_industrial_predicts, - build_industrial, - split_any_industrial, - split_time_series_industrial, - predict_operation_industrial, - predict_industrial, - predict_for_fit_industrial, - update_column_types_industrial, - transform_lagged_industrial, - fit_topo_extractor_industrial, - transform_topo_extractor_industrial, - transform_lagged_for_fit_industrial, - _check_and_correct_window_size_industrial, - transform_smoothing_industrial, - DaskOptunaTuner, - tune_pipeline_industrial, - reproduce_controlled_industrial, - reproduce_industrial] - -DEFAULT_METHODS = [getattr(class_impl[0], class_impl[1]) - for class_impl in FEDOT_METHOD_TO_REPLACE] -DEFAULT_MODELS_TO_REPLACE = [(MODEL_REPO, 'SKLEARN_REG_MODELS'), - (MODEL_REPO, 'SKLEARN_CLF_MODELS'), - (MODEL_REPO, 'FEDOT_PREPROC_MODEL')] +from fedot.core.context import ExecutionContext +from fedot.industrial.industrial_extension import IndustrialExtension def has_no_resample(pipeline: Pipeline): @@ -113,111 +51,85 @@ def has_no_resample(pipeline: Pipeline): """ for node in pipeline.nodes: if node.name == 'resample': - raise ValueError( - f'Pipeline can not have resample operation') + raise ValueError("Pipeline can not have resample operation") return True -class IndustrialModels: - def __init__(self): +def initialize_industrial_context(backend: str = "default") -> ExecutionContext: + context = ExecutionContext() + extension = IndustrialExtension(backend=backend) + extension.apply(context) + return context + +class IndustrialModels: + def __init__(self, backend: str = "default"): self.industrial_data_operation_path = IND_DATA_OPERATION_PATH self.industrial_model_path = IND_MODEL_OPERATION_PATH - self.base_data_operation_path = DEFAULT_DATA_OPERATION_PATH self.base_model_path = DEFAULT_MODEL_OPERATION_PATH - def _replace_operation(self, to_industrial=True, backend: str = 'default'): - method = INDUSTRIAL_REPLACE_METHODS if to_industrial else DEFAULT_METHODS - for class_impl, method_to_replace in zip(FEDOT_METHOD_TO_REPLACE, method): - setattr(class_impl[0], class_impl[1], method_to_replace) - if backend.__contains__('dask'): - model_to_overload = [SKLEARN_REG_MODELS, SKLEARN_CLF_MODELS, FEDOT_PREPROC_MODEL] - overloaded_model = overload_model_implementation(model_to_overload, backend=backend) - for model_impl, new_backend_impl in zip(DEFAULT_MODELS_TO_REPLACE, overloaded_model): - setattr(model_impl[0], model_impl[1], new_backend_impl) - - def setup_repository(self, backend: str = 'default'): - OperationTypesRepository.__repository_dict__.update( - {'data_operation': {'file': self.industrial_data_operation_path, - 'initialized_repo': True, - 'default_tags': []}}) - - OperationTypesRepository.assign_repo( - 'data_operation', self.industrial_data_operation_path) - - OperationTypesRepository.__repository_dict__.update( - {'model': {'file': self.industrial_model_path, - 'initialized_repo': True, - 'default_tags': []}}) - OperationTypesRepository.assign_repo( - 'model', self.industrial_model_path) - # replace mutations - self._replace_operation(to_industrial=True, backend=backend) - - class_rules.append(has_no_data_flow_conflicts_in_industrial_pipeline) - ts_rules.append(has_no_lagged_conflicts_in_ts_pipeline) + self.backend = backend + self.extension = IndustrialExtension(backend=backend) + self.context: ExecutionContext | None = None + + def get_industrial_context(self) -> ExecutionContext: + self.context = ExecutionContext() + self.extension.apply(self.context) + return self.context + + def setup_repository(self) -> OperationTypesRepository: + OperationTypesRepository.__repository_dict__.update({ + 'data_operation': { + 'file': self.industrial_data_operation_path, + 'initialized_repo': True, + 'default_tags': [] + } + }) + OperationTypesRepository.assign_repo('data_operation', self.industrial_data_operation_path) + + OperationTypesRepository.__repository_dict__.update({ + 'model': { + 'file': self.industrial_model_path, + 'initialized_repo': True, + 'default_tags': [] + } + }) + OperationTypesRepository.assign_repo('model', self.industrial_model_path) + + self.get_industrial_context() + self.context.class_rules.append(has_no_data_flow_conflicts_in_industrial_pipeline) + self.context.ts_rules.append(has_no_lagged_conflicts_in_ts_pipeline) + return OperationTypesRepository - def setup_default_repository(self, backend: str = 'default'): - """ - Switching to fedot models. - """ - OperationTypesRepository.__repository_dict__.update( - {'data_operation': {'file': self.base_data_operation_path, - 'initialized_repo': None, - 'default_tags': [ - OperationTypesRepository.DEFAULT_DATA_OPERATION_TAGS]}}) - OperationTypesRepository.assign_repo( - 'data_operation', self.base_data_operation_path) - - OperationTypesRepository.__repository_dict__.update( - {'model': {'file': self.base_model_path, - 'initialized_repo': None, - 'default_tags': []}}) + def setup_default_repository(self) -> OperationTypesRepository: + OperationTypesRepository.__repository_dict__.update({ + 'data_operation': { + 'file': self.base_data_operation_path, + 'initialized_repo': None, + 'default_tags': [OperationTypesRepository.DEFAULT_DATA_OPERATION_TAGS] + } + }) + OperationTypesRepository.assign_repo('data_operation', self.base_data_operation_path) + + OperationTypesRepository.__repository_dict__.update({ + 'model': { + 'file': self.base_model_path, + 'initialized_repo': None, + 'default_tags': [] + } + }) OperationTypesRepository.assign_repo('model', self.base_model_path) - self._replace_operation(to_industrial=False, backend=backend) - common_rules.append(has_no_resample) + self.context = ExecutionContext() + self.context.common_rules.append(has_no_resample) + return OperationTypesRepository - def __enter__(self): - """ - Switching to industrial models - """ - OperationTypesRepository.__repository_dict__.update( - {'data_operation': {'file': self.industrial_data_operation_path, - 'initialized_repo': True, - 'default_tags': []}}) - - OperationTypesRepository.assign_repo( - 'data_operation', self.industrial_data_operation_path) - - OperationTypesRepository.__repository_dict__.update( - {'model': {'file': self.industrial_model_path, - 'initialized_repo': True, - 'default_tags': []}}) - OperationTypesRepository.assign_repo( - 'model', self.industrial_model_path) - - setattr(PipelineSearchSpace, "get_parameters_dict", - get_industrial_search_space) - setattr(ApiComposer, "_get_default_mutations", - _get_default_industrial_mutations) + def __enter__(self) -> ExecutionContext: + self.setup_repository() + return self.context def __exit__(self, exc_type, exc_val, exc_tb): - """ - Switching to fedot models. - """ - OperationTypesRepository.__repository_dict__.update( - {'data_operation': {'file': self.base_data_operation_path, - 'initialized_repo': None, - 'default_tags': [ - OperationTypesRepository.DEFAULT_DATA_OPERATION_TAGS]}}) - OperationTypesRepository.assign_repo( - 'data_operation', self.base_data_operation_path) - - OperationTypesRepository.__repository_dict__.update( - {'model': {'file': self.base_model_path, - 'initialized_repo': None, - 'default_tags': []}}) - OperationTypesRepository.assign_repo('model', self.base_model_path) + self.setup_default_repository() + self.context = None diff --git a/test/integration/pipelines/tuning/test_pipeline_tuning.py b/test/integration/pipelines/tuning/test_pipeline_tuning.py deleted file mode 100644 index 3fcc228243..0000000000 --- a/test/integration/pipelines/tuning/test_pipeline_tuning.py +++ /dev/null @@ -1,565 +0,0 @@ -import os -from time import time - -import pytest -from golem.core.tuning.hyperopt_tuner import get_node_parameters_for_hyperopt -from golem.core.tuning.iopt_tuner import IOptTuner -from golem.core.tuning.optuna_tuner import OptunaTuner -from golem.core.tuning.sequential import SequentialTuner -from golem.core.tuning.simultaneous import SimultaneousTuner -from golem.utilities.data_structures import ensure_wrapped_in_sequence -from hyperopt import hp -from hyperopt.pyll.stochastic import sample as hp_sample - -from examples.simple.time_series_forecasting.ts_pipelines import ts_complex_ridge_smoothing_pipeline, \ - ts_polyfit_ridge_pipeline -from fedot.core.data.data import InputData -from fedot.core.data.data_split import train_test_data_setup -from fedot.core.operations.evaluation.operation_implementations.models.ts_implementations.statsmodels import \ - GLMImplementation -from fedot.core.pipelines.node import PipelineNode -from fedot.core.pipelines.pipeline import Pipeline -from fedot.core.pipelines.pipeline_builder import PipelineBuilder -from fedot.core.pipelines.tuning.search_space import PipelineSearchSpace -from fedot.core.pipelines.tuning.tuner_builder import TunerBuilder -from fedot.core.repository.dataset_types import DataTypesEnum -from fedot.core.repository.metrics_repository import RegressionMetricsEnum, ClassificationMetricsEnum -from fedot.core.repository.tasks import Task, TaskTypesEnum -from fedot.core.utils import fedot_project_root, NESTED_PARAMS_LABEL -from test.unit.multimodal.data_generators import get_single_task_multimodal_tabular_data, get_multimodal_pipeline -from test.unit.tasks.test_forecasting import get_ts_data - - -@pytest.fixture(scope='package') -def regression_dataset(): - test_file_path = str(os.path.dirname(__file__)) - file = os.path.join(str(fedot_project_root()), 'test/data/simple_regression_train.csv') - return InputData.from_csv(os.path.join(test_file_path, file), task=Task(TaskTypesEnum.regression)) - - -@pytest.fixture() -def classification_dataset(): - test_file_path = str(os.path.dirname(__file__)) - file = os.path.join(str(fedot_project_root()), 'test/data/simple_classification.csv') - return InputData.from_csv(os.path.join(test_file_path, file), task=Task(TaskTypesEnum.classification)) - - -@pytest.fixture() -def tiny_classification_dataset(): - test_file_path = str(os.path.dirname(__file__)) - file = os.path.join(str(fedot_project_root()), 'test/data/tiny_simple_classification.csv') - return InputData.from_csv(os.path.join(test_file_path, file), task=Task(TaskTypesEnum.classification)) - - -@pytest.fixture() -def multi_classification_dataset(): - test_file_path = str(os.path.dirname(__file__)) - file = os.path.join(str(fedot_project_root()), 'test/data/multiclass_classification.csv') - return InputData.from_csv(os.path.join(test_file_path, file), task=Task(TaskTypesEnum.classification)) - - -@pytest.fixture() -def ts_forecasting_dataset(): - train_data, _ = get_ts_data(n_steps=700, forecast_length=20) - return train_data - - -@pytest.fixture() -def multimodal_dataset(): - data, _ = get_single_task_multimodal_tabular_data() - return data - - -def get_simple_regr_pipeline(operation_type='rfr'): - final = PipelineNode(operation_type=operation_type) - pipeline = Pipeline(final) - - return pipeline - - -def get_complex_regr_pipeline(): - node_scaling = PipelineNode(operation_type='scaling') - node_ridge = PipelineNode('ridge', nodes_from=[node_scaling]) - node_linear = PipelineNode('linear', nodes_from=[node_scaling]) - final = PipelineNode('rfr', nodes_from=[node_ridge, node_linear]) - pipeline = Pipeline(final) - - return pipeline - - -def get_regr_pipelines(): - simple_pipelines = [get_simple_regr_pipeline(operation_type) for operation_type in get_regr_operation_types()] - - return simple_pipelines + [get_complex_regr_pipeline()] - - -def get_simple_class_pipeline(operation_type='logit'): - final = PipelineNode(operation_type=operation_type) - pipeline = Pipeline(final) - - return pipeline - - -def get_complex_class_pipeline(): - first = PipelineNode(operation_type='knn') - second = PipelineNode(operation_type='pca') - final = PipelineNode(operation_type='logit', - nodes_from=[first, second]) - - pipeline = Pipeline(final) - - return pipeline - - -def get_pipeline_with_no_params_to_tune(): - first = PipelineNode(operation_type='scaling') - final = PipelineNode(operation_type='bernb', - nodes_from=[first]) - - pipeline = Pipeline(final) - - return pipeline - - -def get_class_pipelines(): - simple_pipelines = [get_simple_class_pipeline(operation_type) for operation_type in get_class_operation_types()] - - return simple_pipelines + [get_complex_class_pipeline()] - - -def get_ts_forecasting_pipelines(): - pipelines = [ts_polyfit_ridge_pipeline(2), ts_complex_ridge_smoothing_pipeline()] - return pipelines - - -def get_multimodal_pipelines(): - return [get_multimodal_pipeline()] - - -def get_regr_operation_types(): - return ['lgbmreg'] - - -def get_class_operation_types(): - return ['rf'] - - -def get_regr_losses(): - return [RegressionMetricsEnum.RMSE, RegressionMetricsEnum.MAPE] - - -def get_class_losses(): - return [ClassificationMetricsEnum.ROCAUC, ClassificationMetricsEnum.accuracy] - - -def get_not_default_search_space(): - custom_search_space = { - 'logit': { - 'C': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [1e-1, 5.0], - 'type': 'continuous'} - }, - 'ridge': { - 'alpha': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [0.01, 5.0], - 'type': 'continuous'} - }, - 'lgbmreg': { - 'learning_rate': { - 'hyperopt-dist': hp.loguniform, - 'sampling-scope': [0.03, 0.1], - 'type': 'continuous'}, - 'colsample_bytree': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [0.2, 0.8], - 'type': 'continuous'}, - 'subsample': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [0.1, 0.8], - 'type': 'continuous'} - }, - 'dt': { - 'max_depth': { - 'hyperopt-dist': hp.uniformint, - 'sampling-scope': [1, 5], - 'type': 'discrete'}, - 'min_samples_split': { - 'hyperopt-dist': hp.uniformint, - 'sampling-scope': [10, 25], - 'type': 'discrete'} - }, - 'ar': { - 'lag_1': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [2, 100], - 'type': 'continuous'}, - 'lag_2': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [2, 500], - 'type': 'continuous'} - }, - 'pca': { - 'n_components': { - 'hyperopt-dist': hp.uniform, - 'sampling-scope': [0.1, 0.5], - 'type': 'continuous'} - } - } - return PipelineSearchSpace(custom_search_space=custom_search_space) - - -def run_pipeline_tuner(train_data, - pipeline, - loss_function, - tuner=SimultaneousTuner, - search_space=PipelineSearchSpace(), - cv=None, - iterations=5, - early_stopping_rounds=None, **kwargs): - # if data is time series then lagged window should be tuned correctly - # because lagged window raises error if windows size is uncorrect - # and tuner will fall - if train_data.data_type in (DataTypesEnum.ts, DataTypesEnum.multi_ts): - forecast_length = train_data.task.task_params.forecast_length - folds = cv or 1 - validation_blocks = 1 - max_window = int(train_data.features.shape[0] / (folds + 1)) - (forecast_length * validation_blocks) - 1 - ssp = {'window_size': {'hyperopt-dist': hp.uniformint, 'sampling-scope': [2, max_window], 'type': 'discrete'}} - if search_space.custom_search_space is None: - search_space.custom_search_space = {'lagged': ssp} - else: - search_space.custom_search_space['lagged'] = ssp - search_space.replace_default_search_space = True - search_space.parameters_per_operation = search_space.get_parameters_dict() - - # Pipeline tuning - pipeline_tuner = TunerBuilder(train_data.task) \ - .with_tuner(tuner) \ - .with_metric(loss_function) \ - .with_cv_folds(cv) \ - .with_iterations(iterations) \ - .with_n_jobs(1) \ - .with_early_stopping_rounds(early_stopping_rounds) \ - .with_search_space(search_space) \ - .with_additional_params(**kwargs) \ - .build(train_data) - tuned_pipeline = pipeline_tuner.tune(pipeline, show_progress=False) - return pipeline_tuner, tuned_pipeline - - -def run_node_tuner(train_data, - pipeline, - loss_function, - search_space=PipelineSearchSpace(), - cv=None, - node_index=0, - iterations=3, - early_stopping_rounds=None): - # Pipeline tuning - node_tuner = TunerBuilder(train_data.task) \ - .with_tuner(SequentialTuner) \ - .with_metric(loss_function) \ - .with_cv_folds(cv) \ - .with_iterations(iterations) \ - .with_search_space(search_space) \ - .with_early_stopping_rounds(early_stopping_rounds) \ - .build(train_data) - tuned_pipeline = node_tuner.tune_node(pipeline, node_index) - return node_tuner, tuned_pipeline - - -@pytest.mark.parametrize('data_fixture', ['classification_dataset']) -def test_custom_params_setter(data_fixture, request): - data = request.getfixturevalue(data_fixture) - pipeline = get_complex_class_pipeline() - - custom_params = dict(C=10) - - pipeline.root_node.parameters = custom_params - pipeline.fit(data) - params = pipeline.root_node.fitted_operation.get_params() - - assert params['C'] == 10 - - -@pytest.mark.parametrize('data_fixture, pipelines, loss_functions', - [('regression_dataset', get_regr_pipelines(), get_regr_losses()), - ('classification_dataset', get_class_pipelines(), get_class_losses()), - ('multi_classification_dataset', get_class_pipelines(), get_class_losses()), - ('ts_forecasting_dataset', get_ts_forecasting_pipelines(), get_regr_losses()), - ('multimodal_dataset', get_multimodal_pipelines(), get_class_losses())]) -@pytest.mark.parametrize('tuner', [SimultaneousTuner, SequentialTuner, OptunaTuner]) -def test_pipeline_tuner_correct(data_fixture, pipelines, loss_functions, request, tuner): - """ Test all tuners for pipeline """ - data = request.getfixturevalue(data_fixture) - cvs = [None, 2] - - for pipeline in pipelines: - for loss_function in loss_functions: - for cv in cvs: - print(pipeline) - pipeline_tuner, tuned_pipeline = run_pipeline_tuner(tuner=tuner, - train_data=data, - pipeline=pipeline, - loss_function=loss_function, - cv=cv) - assert pipeline_tuner.obtained_metric is not None - assert tuned_pipeline is not None - assert not tuned_pipeline.is_fitted - - -@pytest.mark.parametrize('tuner', [SimultaneousTuner, SequentialTuner, IOptTuner, OptunaTuner]) -def test_pipeline_tuner_with_no_parameters_to_tune(classification_dataset, tuner): - pipeline = get_pipeline_with_no_params_to_tune() - pipeline_tuner, tuned_pipeline = run_pipeline_tuner(tuner=tuner, - train_data=classification_dataset, - pipeline=pipeline, - loss_function=ClassificationMetricsEnum.ROCAUC, - iterations=20) - assert pipeline_tuner.obtained_metric is not None - assert tuned_pipeline is not None - assert pipeline_tuner.obtained_metric == pipeline_tuner.init_metric - assert not tuned_pipeline.is_fitted - - -@pytest.mark.parametrize('tuner', [SimultaneousTuner, SequentialTuner, OptunaTuner]) -def test_pipeline_tuner_with_initial_params(classification_dataset, tuner): - """ Test all tuners for pipeline with initial parameters """ - # a model - node = PipelineNode(content={'name': 'xgboost', 'params': {'max_depth': 3, - 'learning_rate': 0.03, - 'min_child_weight': 2}}) - pipeline = Pipeline(node) - pipeline_tuner, tuned_pipeline = run_pipeline_tuner(tuner=tuner, - train_data=classification_dataset, - pipeline=pipeline, - loss_function=ClassificationMetricsEnum.ROCAUC, - iterations=20) - assert pipeline_tuner.obtained_metric is not None - assert tuned_pipeline is not None - assert not tuned_pipeline.is_fitted - - -@pytest.mark.parametrize('data_fixture, pipelines, loss_functions', - [('regression_dataset', get_regr_pipelines(), get_regr_losses()), - ('classification_dataset', get_class_pipelines(), get_class_losses()), - ('multi_classification_dataset', get_class_pipelines(), get_class_losses()), - ('ts_forecasting_dataset', get_ts_forecasting_pipelines(), get_regr_losses()), - ('multimodal_dataset', get_multimodal_pipelines(), get_class_losses())]) -@pytest.mark.parametrize('tuner', [SimultaneousTuner, SequentialTuner, OptunaTuner]) -def test_pipeline_tuner_with_custom_search_space(data_fixture, pipelines, loss_functions, request, tuner): - """ Test tuners with different search spaces """ - data = request.getfixturevalue(data_fixture) - train_data, test_data = train_test_data_setup(data=data) - search_spaces = [PipelineSearchSpace(), get_not_default_search_space()] - - for search_space in search_spaces: - pipeline_tuner, tuned_pipeline = run_pipeline_tuner(tuner=tuner, - train_data=train_data, - pipeline=pipelines[0], - loss_function=loss_functions[0], - search_space=search_space) - assert pipeline_tuner.obtained_metric is not None - assert tuned_pipeline is not None - - -@pytest.mark.parametrize('data_fixture, pipelines, loss_functions', - [('regression_dataset', get_regr_pipelines(), get_regr_losses()), - ('classification_dataset', get_class_pipelines(), get_class_losses()), - ('multi_classification_dataset', get_class_pipelines(), get_class_losses()), - ('ts_forecasting_dataset', get_ts_forecasting_pipelines(), get_regr_losses()), - ('multimodal_dataset', get_multimodal_pipelines(), get_class_losses())]) -def test_certain_node_tuning_correct(data_fixture, pipelines, loss_functions, request): - """ Test SequentialTuner for particular node based on hyperopt library """ - data = request.getfixturevalue(data_fixture) - cvs = [None, 2] - - for pipeline in pipelines: - for loss_function in loss_functions: - for cv in cvs: - node_tuner, tuned_pipeline = run_node_tuner(train_data=data, - pipeline=pipeline, - loss_function=loss_function, - cv=cv) - assert node_tuner.obtained_metric is not None - assert not tuned_pipeline.is_fitted - assert tuned_pipeline is not None - - -@pytest.mark.parametrize('data_fixture, pipelines, loss_functions', - [('regression_dataset', get_regr_pipelines(), get_regr_losses()), - ('classification_dataset', get_class_pipelines(), get_class_losses()), - ('multi_classification_dataset', get_class_pipelines(), get_class_losses()), - ('ts_forecasting_dataset', get_ts_forecasting_pipelines(), get_regr_losses()), - ('multimodal_dataset', get_multimodal_pipelines(), get_class_losses())]) -def test_certain_node_tuner_with_custom_search_space(data_fixture, pipelines, loss_functions, request): - """ Test SequentialTuner for particular node with different search spaces """ - data = request.getfixturevalue(data_fixture) - train_data, test_data = train_test_data_setup(data=data) - search_spaces = [PipelineSearchSpace(), get_not_default_search_space()] - - for search_space in search_spaces: - node_tuner, tuned_pipeline = run_node_tuner(train_data=train_data, - pipeline=pipelines[0], - loss_function=loss_functions[0], - search_space=search_space) - assert node_tuner.obtained_metric is not None - assert tuned_pipeline is not None - - -@pytest.mark.parametrize('n_steps', [100, 133, 217, 300]) -@pytest.mark.parametrize('tuner', [SimultaneousTuner, SequentialTuner, IOptTuner, OptunaTuner]) -def test_ts_pipeline_with_stats_model(n_steps, tuner): - """ Tests tuners for time series forecasting task with AR model """ - train_data, test_data = get_ts_data(n_steps=n_steps, forecast_length=5) - - ar_pipeline = Pipeline(PipelineNode('ar')) - - for search_space in [PipelineSearchSpace(), get_not_default_search_space()]: - # Tune AR model - tuner_ar = TunerBuilder(train_data.task) \ - .with_tuner(tuner) \ - .with_metric(RegressionMetricsEnum.MSE) \ - .with_iterations(3) \ - .with_search_space(search_space).build(train_data) - tuned_pipeline = tuner_ar.tune(ar_pipeline, show_progress=False) - assert tuned_pipeline is not None - assert tuner_ar.obtained_metric is not None - - -@pytest.mark.parametrize('data_fixture', ['tiny_classification_dataset']) -def test_early_stop_in_tuning(data_fixture, request): - data = request.getfixturevalue(data_fixture) - train_data, test_data = train_test_data_setup(data=data) - - start_pipeline_tuner = time() - _ = run_pipeline_tuner(tuner=SimultaneousTuner, - train_data=train_data, - pipeline=get_class_pipelines()[0], - loss_function=ClassificationMetricsEnum.ROCAUC, - iterations=1000, - early_stopping_rounds=1) - assert time() - start_pipeline_tuner < 1.3 - - start_sequential_tuner = time() - _ = run_pipeline_tuner(tuner=SequentialTuner, - train_data=train_data, - pipeline=get_class_pipelines()[0], - loss_function=ClassificationMetricsEnum.ROCAUC, - iterations=1000, - early_stopping_rounds=1) - assert time() - start_sequential_tuner < 1.3 - - start_node_tuner = time() - _ = run_node_tuner(train_data=train_data, - pipeline=get_class_pipelines()[0], - loss_function=ClassificationMetricsEnum.ROCAUC, - iterations=1000, - early_stopping_rounds=1) - assert time() - start_node_tuner < 1.3 - - -def test_search_space_correctness_after_customization(): - default_search_space = PipelineSearchSpace() - - custom_search_space = {'gbr': {'max_depth': { - 'hyperopt-dist': hp.choice, - 'sampling-scope': [[3, 7, 31, 127, 8191, 131071]], - 'type': 'categorical'}}} - custom_search_space_without_replace = PipelineSearchSpace(custom_search_space=custom_search_space, - replace_default_search_space=False) - custom_search_space_with_replace = PipelineSearchSpace(custom_search_space=custom_search_space, - replace_default_search_space=True) - - default_params, _ = get_node_parameters_for_hyperopt(default_search_space, - node_id=0, - node=PipelineNode('gbr')) - custom_without_replace_params, _ = get_node_parameters_for_hyperopt(custom_search_space_without_replace, - node_id=0, - node=PipelineNode('gbr')) - custom_with_replace_params, _ = get_node_parameters_for_hyperopt(custom_search_space_with_replace, - node_id=0, - node=PipelineNode('gbr')) - - assert default_params.keys() == custom_without_replace_params.keys() - assert default_params.keys() != custom_with_replace_params.keys() - assert default_params['0 || gbr | max_depth'] != custom_without_replace_params['0 || gbr | max_depth'] - assert default_params['0 || gbr | max_depth'] != custom_with_replace_params['0 || gbr | max_depth'] - - -def test_search_space_get_operation_parameter_range(): - default_search_space = PipelineSearchSpace() - gbr_operations = ['loss', 'learning_rate', 'max_depth', 'min_samples_split', - 'min_samples_leaf', 'subsample', 'max_features', 'alpha'] - - custom_search_space = {'gbr': {'max_depth': { - 'hyperopt-dist': hp.choice, - 'sampling-scope': [[3, 7, 31, 127, 8191, 131071]], - 'type': 'categorical'}}} - custom_search_space_without_replace = PipelineSearchSpace(custom_search_space=custom_search_space, - replace_default_search_space=False) - custom_search_space_with_replace = PipelineSearchSpace(custom_search_space=custom_search_space, - replace_default_search_space=True) - - default_operations = default_search_space.get_parameters_for_operation('gbr') - custom_without_replace_operations = custom_search_space_without_replace.get_parameters_for_operation('gbr') - custom_with_replace_operations = custom_search_space_with_replace.get_parameters_for_operation('gbr') - - assert default_operations == gbr_operations - assert custom_without_replace_operations == gbr_operations - assert custom_with_replace_operations == ['max_depth'] - - -def test_complex_search_space(): - space = PipelineSearchSpace() - for i in range(20): - operation_parameters = space.parameters_per_operation.get("glm") - new_value = hp_sample(operation_parameters[NESTED_PARAMS_LABEL]) - for params in new_value['sampling-scope'][0]: - assert params['link'] in GLMImplementation.family_distribution[params['family']]['available_links'] - - -@pytest.mark.parametrize('tuner', [SimultaneousTuner, SequentialTuner, IOptTuner, OptunaTuner]) -def test_complex_search_space_tuning_correct(tuner): - """ Tests Tuners for time series forecasting task with GLM model that has a complex glm search space""" - train_data, test_data = get_ts_data(n_steps=700, forecast_length=20) - - # ridge added because IOpt requires at least one continuous parameter - glm_pipeline = PipelineBuilder().add_sequence('glm', 'ridge', branch_idx=0).build() - initial_parameters = glm_pipeline.nodes[0].parameters - tuner = TunerBuilder(train_data.task) \ - .with_tuner(tuner) \ - .with_metric(RegressionMetricsEnum.MSE) \ - .with_iterations(100) \ - .build(train_data) - tuned_glm_pipeline = tuner.tune(glm_pipeline) - found_parameters = tuned_glm_pipeline.nodes[0].parameters - assert initial_parameters != found_parameters - - -@pytest.mark.parametrize('data_fixture, pipelines, loss_functions', - [('regression_dataset', get_regr_pipelines(), get_regr_losses()), - ('classification_dataset', get_class_pipelines(), get_class_losses()), - ('multi_classification_dataset', get_class_pipelines(), get_class_losses()), - ('ts_forecasting_dataset', get_ts_forecasting_pipelines(), get_regr_losses()), - ('multimodal_dataset', get_multimodal_pipelines(), get_class_losses())]) -@pytest.mark.parametrize('tuner', [OptunaTuner]) -def test_multiobj_tuning(data_fixture, pipelines, loss_functions, request, tuner): - """ Test multi objective tuning is correct """ - data = request.getfixturevalue(data_fixture) - cvs = [None, 2] - - for pipeline in pipelines: - for cv in cvs: - pipeline_tuner, tuned_pipelines = run_pipeline_tuner(tuner=tuner, - train_data=data, - pipeline=pipeline, - loss_function=loss_functions, - cv=cv) - assert tuned_pipelines is not None - assert all([tuned_pipeline is not None for tuned_pipeline in ensure_wrapped_in_sequence(tuned_pipelines)]) - for metrics in pipeline_tuner.obtained_metric: - assert len(metrics) == len(loss_functions) - assert all(metric is not None for metric in metrics) diff --git a/test/unit/context/__init__.py b/test/unit/context/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/test/unit/data/test_supplementary_data.py b/test/unit/data/test_supplementary_data.py deleted file mode 100644 index 0a4f9beaa1..0000000000 --- a/test/unit/data/test_supplementary_data.py +++ /dev/null @@ -1,128 +0,0 @@ -import numpy as np -import pytest - -from fedot.core.data.data import OutputData -from fedot.core.data.merge.data_merger import DataMerger -from fedot.core.data.merge.supplementary_data_merger import SupplementaryDataMerger -from fedot.core.data.supplementary_data import SupplementaryData -from fedot.core.pipelines.node import PipelineNode -from fedot.core.pipelines.pipeline import Pipeline -from fedot.core.repository.dataset_types import DataTypesEnum -from fedot.core.repository.tasks import Task, TaskTypesEnum -from fedot.preprocessing.data_types import TYPE_TO_ID -from test.unit.data.test_data_merge import unequal_outputs_table # noqa, fixture -from test.unit.tasks.test_regression import get_synthetic_regression_data - - -@pytest.fixture() -def outputs_table_with_different_types(): - """ Create datasets with different types of columns in predictions """ - task = Task(TaskTypesEnum.regression) - idx = [0, 1, 2] - target = [1, 2, 10] - data_info_first = SupplementaryData(col_type_ids={'features': np.array([TYPE_TO_ID[str], TYPE_TO_ID[float]]), - 'target': np.array([TYPE_TO_ID[int]])}) - output_first = OutputData(idx=idx, features=None, - predict=np.array([['a', 1.1], ['b', 2], ['c', 3]], dtype=object), - task=task, target=target, data_type=DataTypesEnum.table, - supplementary_data=data_info_first) - - data_info_second = SupplementaryData(col_type_ids={'features': np.array([TYPE_TO_ID[float]]), - 'target': np.array([TYPE_TO_ID[int]])}) - output_second = OutputData(idx=idx, features=None, - predict=np.array([[2.5], [2.1], [9.3]], dtype=float), - task=task, target=target, data_type=DataTypesEnum.table, - supplementary_data=data_info_second) - - return [output_first, output_second] - - -def generate_straight_pipeline(): - """ Simple linear pipeline """ - node_scaling = PipelineNode('scaling') - node_ridge = PipelineNode('ridge', nodes_from=[node_scaling]) - node_linear = PipelineNode('linear', nodes_from=[node_ridge]) - pipeline = Pipeline(node_linear) - return pipeline - - -def test_parent_mask_correct(unequal_outputs_table): # noqa, fixture - """ Test correctness of function for tables mask generation """ - correct_parent_mask = {'input_ids': [0, 1], 'flow_lens': [1, 0]} - - # Calculate parent mask from outputs - main_output = DataMerger.find_main_output(unequal_outputs_table) - p_mask = SupplementaryDataMerger(unequal_outputs_table, main_output).prepare_parent_mask() - - assert tuple(p_mask['input_ids']) == tuple(correct_parent_mask['input_ids']) - assert tuple(p_mask['flow_lens']) == tuple(correct_parent_mask['flow_lens']) - - -def test_calculate_data_flow_len_correct(): - """ Function checks whether the number of nodes visited by the data block - is calculated correctly """ - - # Pipeline consists of 3 nodes - simple_pipeline = generate_straight_pipeline() - data = get_synthetic_regression_data(n_samples=100, n_features=2) - - simple_pipeline.fit(data) - predict_output = simple_pipeline.predict(data) - - assert predict_output.supplementary_data.data_flow_length == 2 - - -def test_get_compound_mask_correct(): - """ Checking whether the procedure for combining lists with keys is - performed correctly for features_mask """ - - synthetic_mask = {'input_ids': [0, 0, 1, 1], - 'flow_lens': [1, 1, 0, 0]} - output_example = OutputData(idx=[0, 0], features=[0, 0], predict=[0, 0], - task=Task(TaskTypesEnum.regression), - target=[0, 0], data_type=DataTypesEnum.table, - supplementary_data=SupplementaryData(features_mask=synthetic_mask)) - - mask = output_example.supplementary_data.compound_mask - - assert ('01', '01', '10', '10') == tuple(mask) - - -def test_define_parents_with_equal_lengths(): - """ - Check the processing of the case when the decompose operation receives - data whose flow_lens is not different. In this case, the data that came - from the data_operation node is used as the "Data parent". - - Such case is common for time series forecasting pipelines. So we imitate - merged output from ARIMA and lagged operations - """ - sd = SupplementaryData(is_main_target=True, - data_flow_length=1, - features_mask={'input_ids': [0, 0, 0, 1, 1, 1], - 'flow_lens': [0, 0, 0, 0, 0, 0]}, - previous_operations=['arima', 'lagged']) - features_mask = np.array(sd.compound_mask) - unique_features_masks = np.unique(features_mask) - - model_parent, data_parent = sd.define_parents(unique_features_masks, task=TaskTypesEnum.ts_forecasting) - - assert model_parent == '00' - assert data_parent == '10' - - -def test_define_types_after_merging(outputs_table_with_different_types): - """ Check if column types for features table perform correctly """ - outputs = outputs_table_with_different_types - # new_idx, features, target, task, d_type, updated_info = DataMerger(outputs).merge() - merged_data = DataMerger.get(outputs).merge() - updated_info = merged_data.supplementary_data - - feature_type_ids = updated_info.col_type_ids['features'] - target_type_ids = updated_info.col_type_ids['target'] - - # Target type must stay the same - ancestor_target_type = outputs[0].supplementary_data.col_type_ids['target'][0] - assert target_type_ids[0] == ancestor_target_type - assert len(feature_type_ids) == 3 - assert tuple(feature_type_ids) == (TYPE_TO_ID[str], TYPE_TO_ID[float], TYPE_TO_ID[float])