From aa558ed923e4b19973891e6b454fc28f99ea740b Mon Sep 17 00:00:00 2001 From: Vijetha Vijayendran Date: Wed, 15 Mar 2023 13:58:23 -0700 Subject: [PATCH 1/5] Flatten queue_settings properties on sweep and automl jobs --- .../azure/ai/ml/entities/_builders/sweep.py | 8 ++++++ .../ai/ml/entities/_job/automl/automl_job.py | 27 +++++++++++++++++++ .../ai/ml/entities/_job/queue_settings.py | 2 +- .../_job/sweep/parameterized_sweep.py | 20 ++++++++++++++ 4 files changed, 56 insertions(+), 1 deletion(-) diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py index 154f31853ee0..53470b3fc7f7 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py @@ -100,6 +100,10 @@ class Sweep(ParameterizedSweep, BaseNode): UserIdentityConfiguration] :param queue_settings: Queue settings for the job. :type queue_settings: QueueSettings + :param job_tier: Determines the job tier + :type job_tier: str + :param priority: Controls the priority on the compute. + :type priority: str """ def __init__( @@ -125,6 +129,8 @@ def __init__( Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, queue_settings: Optional[QueueSettings] = None, + job_tier: Optional[str] = None, + priority: Optional[str] = None, **kwargs, ): # TODO: get rid of self._job_inputs, self._job_outputs once we have general Input @@ -150,6 +156,8 @@ def __init__( early_termination=early_termination, search_space=search_space, queue_settings=queue_settings, + job_tier=job_tier, + priority=priority, ) self.identity = identity diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py index 78a1345fc973..c4e7b0b16430 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py @@ -51,6 +51,10 @@ class AutoMLJob(Job, JobIOMixin, AutoMLNodeIOMixin, ABC): :paramtype outputs: typing.Optional[Dict[str, str]] :keyword queue_settings: The queue settings for the AutoML job, defaults to None :paramtype queue_settings: typing.Optional[QueueSettings] + :keyword job_tier: Determines the job tier, defaults to None + :paramtype job_tier: typing.Optional[str] + :keyword priority: Controls the priority on the compute, defaults to None + :paramtype priority: typing.Optional[str] :raises ValidationException: task type validation error :raises NotImplementedError: Raises NotImplementedError :return: An AutoML Job @@ -65,6 +69,8 @@ def __init__( Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, queue_settings: Optional[QueueSettings] = None, + job_tier: Optional[str] = None, + priority: Optional[str] = None, **kwargs: Any, ) -> None: """Initialize an AutoML job entity. @@ -84,6 +90,10 @@ def __init__( :paramtype outputs: typing.Optional[Dict[str, str]] :keyword queue_settings: The queue settings for the AutoML job, defaults to None :paramtype queue_settings: typing.Optional[QueueSettings] + :keyword job_tier: Determines the job tier, defaults to None + :paramtype job_tier: typing.Optional[str] + :keyword priority: Controls the priority on the compute, defaults to None + :paramtype priority: typing.Optional[str] :raises ValidationException: task type validation error :raises NotImplementedError: Raises NotImplementedError :return: An AutoML Job @@ -99,6 +109,8 @@ def __init__( self.resources = resources self.identity = identity self.queue_settings = queue_settings + if job_tier is not None or priority is not None: + self.set_queue_settings(job_tier=job_tier, priority=priority) @property @abstractmethod @@ -132,6 +144,21 @@ def test_data(self) -> Input: :rtype: Input """ raise NotImplementedError() + + def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Optional[str] = None): + """Set QueueSettings for the job. + :param job_tier: determines the job tier. + :type job_tier: str + :param priority: controls the priority on the compute. + :type priority: str + """ + if self.queue_settings is None: + self.queue_settings = QueueSettings + + if job_tier is not None: + self.queue_settings.job_tier = job_tier + if priority is not None: + self.queue_settings.priority = priority @classmethod def _load_from_rest(cls, obj: JobBase) -> "AutoMLJob": diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/queue_settings.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/queue_settings.py index ebb85ae6532d..94ea4b1b563a 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/queue_settings.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/queue_settings.py @@ -24,7 +24,7 @@ class QueueSettings(RestTranslatableMixin, DictMixin): "Standard", "Premium". :vartype job_tier: str or ~azure.mgmt.machinelearningservices.models.JobTier :ivar priority: Controls the priority of the job on a compute. - :vartype priority: int + :vartype priority: str """ def __init__( diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py index d1b9e385c00b..0e91573dab0e 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py @@ -72,6 +72,8 @@ def __init__( ] ] = None, queue_settings: Optional[QueueSettings] = None, + job_tier: Optional[str] = None, + priority: Optional[str] = None, ): self.sampling_algorithm = sampling_algorithm self.early_termination = early_termination @@ -84,6 +86,9 @@ def __init__( else: self.objective = objective + if job_tier is not None or priority is not None: + self.set_queue_settings(job_tier=job_tier, priority=priority) + @property def limits(self): return self._limits @@ -137,6 +142,21 @@ def set_limits( if trial_timeout is not None: self.limits.trial_timeout = trial_timeout + def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Optional[str] = None): + """Set QueueSettings for the job. + :param job_tier: determines the job tier. + :type job_tier: str + :param priority: controls the priority on the compute. + :type priority: str + """ + if self.queue_settings is None: + self.queue_settings = QueueSettings + + if job_tier is not None: + self.queue_settings.job_tier = job_tier + if priority is not None: + self.queue_settings.priority = priority + def set_objective(self, *, goal: Optional[str] = None, primary_metric: Optional[str] = None) -> None: """Set the sweep object. From 5b443801ad2c5309fc0d030487755fa347fa68ba Mon Sep 17 00:00:00 2001 From: Vijetha Vijayendran Date: Wed, 15 Mar 2023 14:33:45 -0700 Subject: [PATCH 2/5] Fix typo --- .../azure/ai/ml/entities/_job/automl/automl_job.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py index c4e7b0b16430..4002a517c559 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py @@ -11,7 +11,6 @@ from azure.ai.ml._restclient.v2023_02_01_preview.models import ( JobBase, MLTableJobInput, - QueueSettings, ResourceConfiguration, TaskType, ) @@ -28,6 +27,7 @@ from azure.ai.ml.entities._job.job import Job from azure.ai.ml.entities._job.job_io_mixin import JobIOMixin from azure.ai.ml.entities._job.pipeline._io import AutoMLNodeIOMixin +from azure.ai.ml.entities._job.queue_settings import QueueSettings from azure.ai.ml.exceptions import ErrorCategory, ErrorTarget, ValidationException module_logger = logging.getLogger(__name__) @@ -153,7 +153,7 @@ def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Option :type priority: str """ if self.queue_settings is None: - self.queue_settings = QueueSettings + self.queue_settings = QueueSettings() if job_tier is not None: self.queue_settings.job_tier = job_tier From e5f09eeb3d1d97d9896362fbbfbc0ae194ad03a3 Mon Sep 17 00:00:00 2001 From: Vijetha Vijayendran Date: Wed, 15 Mar 2023 14:40:06 -0700 Subject: [PATCH 3/5] Fix typo --- .../azure/ai/ml/entities/_job/sweep/parameterized_sweep.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py index 0e91573dab0e..a8de9cb532a9 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py @@ -150,7 +150,7 @@ def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Option :type priority: str """ if self.queue_settings is None: - self.queue_settings = QueueSettings + self.queue_settings = QueueSettings() if job_tier is not None: self.queue_settings.job_tier = job_tier From 72aba810397bca459f5b0cb055438fc09bcf818c Mon Sep 17 00:00:00 2001 From: Vijetha Vijayendran Date: Thu, 16 Mar 2023 11:58:37 -0700 Subject: [PATCH 4/5] Add flattened property to sweep on command builder --- .../azure/ai/ml/entities/_builders/command.py | 21 ++- .../azure/ai/ml/entities/_builders/sweep.py | 14 -- .../ai/ml/entities/_job/automl/automl_job.py | 143 ++---------------- .../_job/sweep/parameterized_sweep.py | 23 --- 4 files changed, 29 insertions(+), 172 deletions(-) diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/command.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/command.py index 47da960a0f36..9dd4dd0a05e0 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/command.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/command.py @@ -392,6 +392,12 @@ def set_limits(self, *, timeout: int, **kwargs): # pylint: disable=unused-argum self.limits = CommandJobLimits(timeout=timeout) def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Optional[str] = None): + """Set QueueSettings for the job. + :param job_tier: determines the job tier. + :type job_tier: str + :param priority: controls the priority on the compute. + :type priority: str + """ if isinstance(self.queue_settings, QueueSettings): self.queue_settings.job_tier = job_tier self.queue_settings.priority = priority @@ -422,6 +428,8 @@ def sweep( Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, queue_settings: Optional[QueueSettings] = None, + job_tier: Optional[str] = None, + priority: Optional[str] = None, ) -> Sweep: """Turn the command into a sweep node with extra sweep run setting. The command component in current Command node will be used as its trial component. A command node can sweep for multiple times, and the generated sweep @@ -456,6 +464,10 @@ def sweep( UserIdentityConfiguration] :param queue_settings: Queue settings for the job. :type queue_settings: QueueSettings + :param job_tier: determines the job tier. + :type job_tier: str + :param priority: controls the priority on the compute. + :type priority: str :return: A sweep node with component from current Command node as its trial component. :rtype: Sweep """ @@ -466,6 +478,13 @@ def sweep( if search_space: inputs_search_space.update(search_space) + if not queue_settings: + queue_settings = self.queue_settings + if job_tier is not None: + queue_settings.job_tier = job_tier + if priority is not None: + queue_settings.priority = priority + sweep_node = Sweep( trial=copy.deepcopy( self.component @@ -485,7 +504,7 @@ def sweep( experiment_name=self.experiment_name, identity=self.identity if not identity else identity, _from_component_func=True, - queue_settings=self.queue_settings if queue_settings is None else queue_settings, + queue_settings=queue_settings, ) sweep_node.set_limits( max_total_trials=max_total_trials, diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py index 53470b3fc7f7..3d76a039a58e 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py @@ -22,7 +22,6 @@ from azure.ai.ml.entities._inputs_outputs import Input, Output from azure.ai.ml.entities._job.job_limits import SweepJobLimits from azure.ai.ml.entities._job.pipeline._io import NodeInput -from azure.ai.ml.entities._job.queue_settings import QueueSettings from azure.ai.ml.entities._job.sweep.early_termination_policy import ( BanditPolicy, EarlyTerminationPolicy, @@ -98,12 +97,6 @@ class Sweep(ParameterizedSweep, BaseNode): ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] - :param queue_settings: Queue settings for the job. - :type queue_settings: QueueSettings - :param job_tier: Determines the job tier - :type job_tier: str - :param priority: Controls the priority on the compute. - :type priority: str """ def __init__( @@ -128,9 +121,6 @@ def __init__( identity: Optional[ Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, - queue_settings: Optional[QueueSettings] = None, - job_tier: Optional[str] = None, - priority: Optional[str] = None, **kwargs, ): # TODO: get rid of self._job_inputs, self._job_outputs once we have general Input @@ -155,9 +145,6 @@ def __init__( limits=limits, early_termination=early_termination, search_space=search_space, - queue_settings=queue_settings, - job_tier=job_tier, - priority=priority, ) self.identity = identity @@ -306,7 +293,6 @@ def _to_job(self) -> SweepJob: inputs=self._job_inputs, outputs=self._job_outputs, identity=self.identity, - queue_settings=self.queue_settings, ) @classmethod diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py index 4002a517c559..30c8a9ec042d 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py @@ -8,12 +8,7 @@ from abc import ABC, abstractmethod from typing import Any, Dict, Optional, Union -from azure.ai.ml._restclient.v2023_02_01_preview.models import ( - JobBase, - MLTableJobInput, - ResourceConfiguration, - TaskType, -) +from azure.ai.ml._restclient.v2023_02_01_preview.models import JobBase, MLTableJobInput, ResourceConfiguration, TaskType from azure.ai.ml._utils.utils import camel_to_snake from azure.ai.ml.constants import JobType from azure.ai.ml.constants._common import TYPE, AssetTypes @@ -27,39 +22,13 @@ from azure.ai.ml.entities._job.job import Job from azure.ai.ml.entities._job.job_io_mixin import JobIOMixin from azure.ai.ml.entities._job.pipeline._io import AutoMLNodeIOMixin -from azure.ai.ml.entities._job.queue_settings import QueueSettings from azure.ai.ml.exceptions import ErrorCategory, ErrorTarget, ValidationException module_logger = logging.getLogger(__name__) class AutoMLJob(Job, JobIOMixin, AutoMLNodeIOMixin, ABC): - """Initialize an AutoML job entity. - - Constructor for an AutoMLJob. - - :keyword resources: Resource configuration for the AutoML job, defaults to None - :paramtype resources: typing.Optional[ResourceConfiguration] - :keyword identity: Identity that training job will use while running on compute, defaults to None - :paramtype identity: typing.Optional[ typing.Union[ManagedIdentityConfiguration, AmlTokenConfiguration - , UserIdentityConfiguration] ] - :keyword environment_id: The environment id for the AutoML job, defaults to None - :paramtype environment_id: typing.Optional[str] - :keyword environment_variables: The environment variables for the AutoML job, defaults to None - :paramtype environment_variables: typing.Optional[Dict[str, str]] - :keyword outputs: The outputs for the AutoML job, defaults to None - :paramtype outputs: typing.Optional[Dict[str, str]] - :keyword queue_settings: The queue settings for the AutoML job, defaults to None - :paramtype queue_settings: typing.Optional[QueueSettings] - :keyword job_tier: Determines the job tier, defaults to None - :paramtype job_tier: typing.Optional[str] - :keyword priority: Controls the priority on the compute, defaults to None - :paramtype priority: typing.Optional[str] - :raises ValidationException: task type validation error - :raises NotImplementedError: Raises NotImplementedError - :return: An AutoML Job - :rtype: AutoMLJob - """ + """AutoML job entity.""" def __init__( self, @@ -68,36 +37,15 @@ def __init__( identity: Optional[ Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, - queue_settings: Optional[QueueSettings] = None, - job_tier: Optional[str] = None, - priority: Optional[str] = None, **kwargs: Any, ) -> None: """Initialize an AutoML job entity. - Constructor for an AutoMLJob. - - :keyword resources: Resource configuration for the AutoML job, defaults to None - :paramtype resources: typing.Optional[ResourceConfiguration] - :keyword identity: Identity that training job will use while running on compute, defaults to None - :paramtype identity: typing.Optional[ typing.Union[ManagedIdentityConfiguration, AmlTokenConfiguration - , UserIdentityConfiguration] ] - :keyword environment_id: The environment id for the AutoML job, defaults to None - :paramtype environment_id: typing.Optional[str] - :keyword environment_variables: The environment variables for the AutoML job, defaults to None - :paramtype environment_variables: typing.Optional[Dict[str, str]] - :keyword outputs: The outputs for the AutoML job, defaults to None - :paramtype outputs: typing.Optional[Dict[str, str]] - :keyword queue_settings: The queue settings for the AutoML job, defaults to None - :paramtype queue_settings: typing.Optional[QueueSettings] - :keyword job_tier: Determines the job tier, defaults to None - :paramtype job_tier: typing.Optional[str] - :keyword priority: Controls the priority on the compute, defaults to None - :paramtype priority: typing.Optional[str] - :raises ValidationException: task type validation error - :raises NotImplementedError: Raises NotImplementedError - :return: An AutoML Job - :rtype: AutoMLJob + :param task_details: The task configuration of the job. This can be Classification, Regression, etc. + :param resources: Resource configuration for the job. + :param identity: Identity that training job will use while running on compute. + :type identity: Union[ManagedIdentity, AmlToken, UserIdentity] + :param kwargs: """ kwargs[TYPE] = JobType.AUTOML self.environment_id = kwargs.pop("environment_id", None) @@ -108,68 +56,24 @@ def __init__( self.resources = resources self.identity = identity - self.queue_settings = queue_settings - if job_tier is not None or priority is not None: - self.set_queue_settings(job_tier=job_tier, priority=priority) @property @abstractmethod def training_data(self) -> Input: - """The training data for the AutoML job. - - :raises NotImplementedError: Raises NotImplementedError - :return: Returns the training data for the AutoML job. - :rtype: Input - """ raise NotImplementedError() @property @abstractmethod def validation_data(self) -> Input: - """The validation data for the AutoML job. - - :raises NotImplementedError: Raises NotImplementedError - :return: Returns the validation data for the AutoML job. - :rtype: Input - """ raise NotImplementedError() @property @abstractmethod def test_data(self) -> Input: - """The test data for the AutoML job. - - :raises NotImplementedError: Raises NotImplementedError - :return: Returns the test data for the AutoML job. - :rtype: Input - """ raise NotImplementedError() - - def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Optional[str] = None): - """Set QueueSettings for the job. - :param job_tier: determines the job tier. - :type job_tier: str - :param priority: controls the priority on the compute. - :type priority: str - """ - if self.queue_settings is None: - self.queue_settings = QueueSettings() - - if job_tier is not None: - self.queue_settings.job_tier = job_tier - if priority is not None: - self.queue_settings.priority = priority @classmethod def _load_from_rest(cls, obj: JobBase) -> "AutoMLJob": - """Loads the rest object to a dict containing items to init the AutoMLJob objects. - - :param obj: Azure Resource Manager resource envelope. - :type obj: JobBase - :raises ValidationException: task type validation error - :return: An AutoML Job - :rtype: AutoMLJob - """ task_type = ( camel_to_snake(obj.properties.task_details.task_type) if obj.properties.task_details.task_type else None ) @@ -193,19 +97,6 @@ def _load_from_dict( additional_message: str, **kwargs, ) -> "AutoMLJob": - """Loads the dictionary objects to an AutoMLJob object. - - :param data: A data dictionary. - :type data: typing.Dict - :param context: A context dictionary. - :type context: typing.Dict - :param additional_message: An additional message to be logged in the ValidationException. - :type additional_message: str - - :raises ValidationException: task type validation error - :return: An AutoML Job - :rtype: AutoMLJob - """ task_type = data.get(AutoMLConstants.TASK_TYPE_YAML) class_type = cls._get_task_mapping().get(task_type, None) if class_type: @@ -225,14 +116,7 @@ def _load_from_dict( @classmethod def _create_instance_from_schema_dict(cls, loaded_data: Dict) -> "AutoMLJob": - """Create an automl job instance from schema parsed dict. - - :param loaded_data: A loaded_data dictionary. - :type loaded_data: typing.Dict - :raises ValidationException: task type validation error - :return: An AutoML Job - :rtype: AutoMLJob - """ + """Create an automl job instance from schema parsed dict.""" task_type = loaded_data.pop(AutoMLConstants.TASK_TYPE_YAML) class_type = cls._get_task_mapping().get(task_type, None) if class_type: @@ -247,11 +131,6 @@ def _create_instance_from_schema_dict(cls, loaded_data: Dict) -> "AutoMLJob": @classmethod def _get_task_mapping(cls): - """Create a mapping of task type to job class. - - :return: An AutoMLVertical object containing the task type to job class mapping. - :rtype: AutoMLVertical - """ from .image import ( ImageClassificationJob, ImageClassificationMultilabelJob, @@ -276,11 +155,7 @@ def _get_task_mapping(cls): } def _resolve_data_inputs(self, rest_job): # pylint: disable=no-self-use - """Resolve JobInputs to MLTableJobInputs within data_settings. - - :param rest_job: The rest job object. - :type rest_job: AutoMLJob - """ + """Resolve JobInputs to MLTableJobInputs within data_settings.""" if isinstance(rest_job.training_data, Input): rest_job.training_data = MLTableJobInput(uri=rest_job.training_data.path) if isinstance(rest_job.validation_data, Input): diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py index a8de9cb532a9..4c12f15e8ac1 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py @@ -6,7 +6,6 @@ from azure.ai.ml.exceptions import ErrorCategory, ErrorTarget, ValidationErrorType, ValidationException from ..job_limits import SweepJobLimits -from ..queue_settings import QueueSettings from .early_termination_policy import ( BanditPolicy, EarlyTerminationPolicy, @@ -71,24 +70,17 @@ def __init__( ], ] ] = None, - queue_settings: Optional[QueueSettings] = None, - job_tier: Optional[str] = None, - priority: Optional[str] = None, ): self.sampling_algorithm = sampling_algorithm self.early_termination = early_termination self._limits = limits self.search_space = search_space - self.queue_settings = queue_settings if isinstance(objective, Dict): self.objective = Objective(**objective) else: self.objective = objective - if job_tier is not None or priority is not None: - self.set_queue_settings(job_tier=job_tier, priority=priority) - @property def limits(self): return self._limits @@ -142,21 +134,6 @@ def set_limits( if trial_timeout is not None: self.limits.trial_timeout = trial_timeout - def set_queue_settings(self, *, job_tier: Optional[str] = None, priority: Optional[str] = None): - """Set QueueSettings for the job. - :param job_tier: determines the job tier. - :type job_tier: str - :param priority: controls the priority on the compute. - :type priority: str - """ - if self.queue_settings is None: - self.queue_settings = QueueSettings() - - if job_tier is not None: - self.queue_settings.job_tier = job_tier - if priority is not None: - self.queue_settings.priority = priority - def set_objective(self, *, goal: Optional[str] = None, primary_metric: Optional[str] = None) -> None: """Set the sweep object. From eee37446953ef0b7cb588ec5815401ff48a5804b Mon Sep 17 00:00:00 2001 From: Vijetha Vijayendran Date: Thu, 16 Mar 2023 12:03:36 -0700 Subject: [PATCH 5/5] Fix code merge --- .../azure/ai/ml/entities/_builders/sweep.py | 6 + .../ai/ml/entities/_job/automl/automl_job.py | 116 ++++++++++++++++-- .../_job/sweep/parameterized_sweep.py | 3 + 3 files changed, 116 insertions(+), 9 deletions(-) diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py index 3d76a039a58e..154f31853ee0 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_builders/sweep.py @@ -22,6 +22,7 @@ from azure.ai.ml.entities._inputs_outputs import Input, Output from azure.ai.ml.entities._job.job_limits import SweepJobLimits from azure.ai.ml.entities._job.pipeline._io import NodeInput +from azure.ai.ml.entities._job.queue_settings import QueueSettings from azure.ai.ml.entities._job.sweep.early_termination_policy import ( BanditPolicy, EarlyTerminationPolicy, @@ -97,6 +98,8 @@ class Sweep(ParameterizedSweep, BaseNode): ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] + :param queue_settings: Queue settings for the job. + :type queue_settings: QueueSettings """ def __init__( @@ -121,6 +124,7 @@ def __init__( identity: Optional[ Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, + queue_settings: Optional[QueueSettings] = None, **kwargs, ): # TODO: get rid of self._job_inputs, self._job_outputs once we have general Input @@ -145,6 +149,7 @@ def __init__( limits=limits, early_termination=early_termination, search_space=search_space, + queue_settings=queue_settings, ) self.identity = identity @@ -293,6 +298,7 @@ def _to_job(self) -> SweepJob: inputs=self._job_inputs, outputs=self._job_outputs, identity=self.identity, + queue_settings=self.queue_settings, ) @classmethod diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py index 30c8a9ec042d..78a1345fc973 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/automl/automl_job.py @@ -8,7 +8,13 @@ from abc import ABC, abstractmethod from typing import Any, Dict, Optional, Union -from azure.ai.ml._restclient.v2023_02_01_preview.models import JobBase, MLTableJobInput, ResourceConfiguration, TaskType +from azure.ai.ml._restclient.v2023_02_01_preview.models import ( + JobBase, + MLTableJobInput, + QueueSettings, + ResourceConfiguration, + TaskType, +) from azure.ai.ml._utils.utils import camel_to_snake from azure.ai.ml.constants import JobType from azure.ai.ml.constants._common import TYPE, AssetTypes @@ -28,7 +34,28 @@ class AutoMLJob(Job, JobIOMixin, AutoMLNodeIOMixin, ABC): - """AutoML job entity.""" + """Initialize an AutoML job entity. + + Constructor for an AutoMLJob. + + :keyword resources: Resource configuration for the AutoML job, defaults to None + :paramtype resources: typing.Optional[ResourceConfiguration] + :keyword identity: Identity that training job will use while running on compute, defaults to None + :paramtype identity: typing.Optional[ typing.Union[ManagedIdentityConfiguration, AmlTokenConfiguration + , UserIdentityConfiguration] ] + :keyword environment_id: The environment id for the AutoML job, defaults to None + :paramtype environment_id: typing.Optional[str] + :keyword environment_variables: The environment variables for the AutoML job, defaults to None + :paramtype environment_variables: typing.Optional[Dict[str, str]] + :keyword outputs: The outputs for the AutoML job, defaults to None + :paramtype outputs: typing.Optional[Dict[str, str]] + :keyword queue_settings: The queue settings for the AutoML job, defaults to None + :paramtype queue_settings: typing.Optional[QueueSettings] + :raises ValidationException: task type validation error + :raises NotImplementedError: Raises NotImplementedError + :return: An AutoML Job + :rtype: AutoMLJob + """ def __init__( self, @@ -37,15 +64,30 @@ def __init__( identity: Optional[ Union[ManagedIdentityConfiguration, AmlTokenConfiguration, UserIdentityConfiguration] ] = None, + queue_settings: Optional[QueueSettings] = None, **kwargs: Any, ) -> None: """Initialize an AutoML job entity. - :param task_details: The task configuration of the job. This can be Classification, Regression, etc. - :param resources: Resource configuration for the job. - :param identity: Identity that training job will use while running on compute. - :type identity: Union[ManagedIdentity, AmlToken, UserIdentity] - :param kwargs: + Constructor for an AutoMLJob. + + :keyword resources: Resource configuration for the AutoML job, defaults to None + :paramtype resources: typing.Optional[ResourceConfiguration] + :keyword identity: Identity that training job will use while running on compute, defaults to None + :paramtype identity: typing.Optional[ typing.Union[ManagedIdentityConfiguration, AmlTokenConfiguration + , UserIdentityConfiguration] ] + :keyword environment_id: The environment id for the AutoML job, defaults to None + :paramtype environment_id: typing.Optional[str] + :keyword environment_variables: The environment variables for the AutoML job, defaults to None + :paramtype environment_variables: typing.Optional[Dict[str, str]] + :keyword outputs: The outputs for the AutoML job, defaults to None + :paramtype outputs: typing.Optional[Dict[str, str]] + :keyword queue_settings: The queue settings for the AutoML job, defaults to None + :paramtype queue_settings: typing.Optional[QueueSettings] + :raises ValidationException: task type validation error + :raises NotImplementedError: Raises NotImplementedError + :return: An AutoML Job + :rtype: AutoMLJob """ kwargs[TYPE] = JobType.AUTOML self.environment_id = kwargs.pop("environment_id", None) @@ -56,24 +98,51 @@ def __init__( self.resources = resources self.identity = identity + self.queue_settings = queue_settings @property @abstractmethod def training_data(self) -> Input: + """The training data for the AutoML job. + + :raises NotImplementedError: Raises NotImplementedError + :return: Returns the training data for the AutoML job. + :rtype: Input + """ raise NotImplementedError() @property @abstractmethod def validation_data(self) -> Input: + """The validation data for the AutoML job. + + :raises NotImplementedError: Raises NotImplementedError + :return: Returns the validation data for the AutoML job. + :rtype: Input + """ raise NotImplementedError() @property @abstractmethod def test_data(self) -> Input: + """The test data for the AutoML job. + + :raises NotImplementedError: Raises NotImplementedError + :return: Returns the test data for the AutoML job. + :rtype: Input + """ raise NotImplementedError() @classmethod def _load_from_rest(cls, obj: JobBase) -> "AutoMLJob": + """Loads the rest object to a dict containing items to init the AutoMLJob objects. + + :param obj: Azure Resource Manager resource envelope. + :type obj: JobBase + :raises ValidationException: task type validation error + :return: An AutoML Job + :rtype: AutoMLJob + """ task_type = ( camel_to_snake(obj.properties.task_details.task_type) if obj.properties.task_details.task_type else None ) @@ -97,6 +166,19 @@ def _load_from_dict( additional_message: str, **kwargs, ) -> "AutoMLJob": + """Loads the dictionary objects to an AutoMLJob object. + + :param data: A data dictionary. + :type data: typing.Dict + :param context: A context dictionary. + :type context: typing.Dict + :param additional_message: An additional message to be logged in the ValidationException. + :type additional_message: str + + :raises ValidationException: task type validation error + :return: An AutoML Job + :rtype: AutoMLJob + """ task_type = data.get(AutoMLConstants.TASK_TYPE_YAML) class_type = cls._get_task_mapping().get(task_type, None) if class_type: @@ -116,7 +198,14 @@ def _load_from_dict( @classmethod def _create_instance_from_schema_dict(cls, loaded_data: Dict) -> "AutoMLJob": - """Create an automl job instance from schema parsed dict.""" + """Create an automl job instance from schema parsed dict. + + :param loaded_data: A loaded_data dictionary. + :type loaded_data: typing.Dict + :raises ValidationException: task type validation error + :return: An AutoML Job + :rtype: AutoMLJob + """ task_type = loaded_data.pop(AutoMLConstants.TASK_TYPE_YAML) class_type = cls._get_task_mapping().get(task_type, None) if class_type: @@ -131,6 +220,11 @@ def _create_instance_from_schema_dict(cls, loaded_data: Dict) -> "AutoMLJob": @classmethod def _get_task_mapping(cls): + """Create a mapping of task type to job class. + + :return: An AutoMLVertical object containing the task type to job class mapping. + :rtype: AutoMLVertical + """ from .image import ( ImageClassificationJob, ImageClassificationMultilabelJob, @@ -155,7 +249,11 @@ def _get_task_mapping(cls): } def _resolve_data_inputs(self, rest_job): # pylint: disable=no-self-use - """Resolve JobInputs to MLTableJobInputs within data_settings.""" + """Resolve JobInputs to MLTableJobInputs within data_settings. + + :param rest_job: The rest job object. + :type rest_job: AutoMLJob + """ if isinstance(rest_job.training_data, Input): rest_job.training_data = MLTableJobInput(uri=rest_job.training_data.path) if isinstance(rest_job.validation_data, Input): diff --git a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py index 4c12f15e8ac1..d1b9e385c00b 100644 --- a/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py +++ b/sdk/ml/azure-ai-ml/azure/ai/ml/entities/_job/sweep/parameterized_sweep.py @@ -6,6 +6,7 @@ from azure.ai.ml.exceptions import ErrorCategory, ErrorTarget, ValidationErrorType, ValidationException from ..job_limits import SweepJobLimits +from ..queue_settings import QueueSettings from .early_termination_policy import ( BanditPolicy, EarlyTerminationPolicy, @@ -70,11 +71,13 @@ def __init__( ], ] ] = None, + queue_settings: Optional[QueueSettings] = None, ): self.sampling_algorithm = sampling_algorithm self.early_termination = early_termination self._limits = limits self.search_space = search_space + self.queue_settings = queue_settings if isinstance(objective, Dict): self.objective = Objective(**objective)