diff --git a/kubeflow/spark/api/spark_client.py b/kubeflow/spark/api/spark_client.py index 332eaa286..9f8b84dc0 100644 --- a/kubeflow/spark/api/spark_client.py +++ b/kubeflow/spark/api/spark_client.py @@ -27,6 +27,7 @@ Executor, FileJob, FuncJob, + ResourceDict, SparkConnectInfo, SparkJob, SparkJobStatus, @@ -64,7 +65,7 @@ def connect( base_url: str | None = None, token: str | None = None, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, spark_conf: dict[str, str] | None = None, driver: Driver | None = None, executor: Executor | None = None, @@ -83,8 +84,8 @@ def connect( If provided, connects to existing server. If None, creates new session. token: Optional authentication token for existing server. num_executors: Number of executor instances (create mode only). - resources_per_executor: Resource requirements per executor as dict. - Format: `{"cpu": "5", "memory": "10Gi"}` (create mode only). + resources_per_executor: Resource requirements per executor as ResourceDict. + Format: `{"cpu": 5, "memory": "10Gi"}` (create mode only). spark_conf: Spark configuration dictionary (create mode only). driver: Driver configuration object (create mode only). executor: Executor configuration object (create mode only). @@ -175,7 +176,7 @@ def submit_job( self, job: FileJob | FuncJob, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, spark_conf: dict[str, str] | None = None, options: list | None = None, ) -> str: @@ -196,8 +197,8 @@ def submit_job( Number of executor instances. resources_per_executor: - Resource requirements per executor. - Format: ``{"cpu": "5", "memory": "10Gi"}``. + Resource requirements per executor as ResourceDict. + Format: ``{"cpu": 5, "memory": "10Gi"}``. spark_conf: Spark configuration properties. diff --git a/kubeflow/spark/backends/base.py b/kubeflow/spark/backends/base.py index 0f41a18a6..58f06d591 100644 --- a/kubeflow/spark/backends/base.py +++ b/kubeflow/spark/backends/base.py @@ -24,6 +24,7 @@ Executor, FileJob, FuncJob, + ResourceDict, SparkConnectInfo, SparkJob, SparkJobStatus, @@ -44,7 +45,7 @@ class RuntimeBackend(abc.ABC): def create_and_connect( self, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, spark_conf: dict[str, str] | None = None, driver: Driver | None = None, executor: Executor | None = None, @@ -59,7 +60,7 @@ def create_and_connect( Args: num_executors: Number of executor instances. - resources_per_executor: Resource requirements per executor. + resources_per_executor: Resource requirements per executor as ResourceDict. spark_conf: Spark configuration properties. driver: Driver configuration. executor: Executor configuration. @@ -148,7 +149,7 @@ def submit_job( self, job: FileJob | FuncJob, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, options: list | None = None, spark_conf: dict[str, str] | None = None, ) -> SparkJob: @@ -157,7 +158,7 @@ def submit_job( Args: job: Spark application to execute. num_executors: Number of executor instances. - resources_per_executor: Resource requirements for each executor. + resources_per_executor: Resource requirements for each executor as ResourceDict. options: List of configuration options. Use the Name option for a custom job name. spark_conf: Spark configuration properties to set on the SparkApplication. diff --git a/kubeflow/spark/backends/kubernetes/backend.py b/kubeflow/spark/backends/kubernetes/backend.py index 6d75874c8..aaf566581 100644 --- a/kubeflow/spark/backends/kubernetes/backend.py +++ b/kubeflow/spark/backends/kubernetes/backend.py @@ -55,6 +55,7 @@ Executor, FileJob, FuncJob, + ResourceDict, SparkConnectInfo, SparkConnectState, SparkJob, @@ -121,7 +122,7 @@ def __init__(self, backend_config: KubernetesBackendConfig): def _create_session( self, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, spark_conf: dict[str, str] | None = None, driver: Driver | None = None, executor: Executor | None = None, @@ -131,7 +132,7 @@ def _create_session( Args: num_executors: Number of executor instances. - resources_per_executor: Resource requirements per executor. + resources_per_executor: Resource requirements per executor as ResourceDict. spark_conf: Spark configuration properties. driver: Driver configuration. executor: Executor configuration. @@ -628,7 +629,7 @@ def _get_or_create() -> None: def create_and_connect( self, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, spark_conf: dict[str, str] | None = None, driver: Driver | None = None, executor: Executor | None = None, @@ -645,7 +646,7 @@ def create_and_connect( Args: num_executors: Number of executor instances. - resources_per_executor: Resource requirements per executor. + resources_per_executor: Resource requirements per executor as ResourceDict. spark_conf: Spark configuration properties. driver: Driver configuration. executor: Executor configuration. @@ -864,7 +865,7 @@ def submit_job( self, job: FileJob | FuncJob, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, options: list | None = None, spark_conf: dict[str, str] | None = None, ) -> SparkJob: @@ -878,7 +879,7 @@ def submit_job( Number of executor instances. resources_per_executor: - Resource requirements per executor. + Resource requirements per executor as ResourceDict. options: List of additional Spark configuration options. diff --git a/kubeflow/spark/backends/kubernetes/utils.py b/kubeflow/spark/backends/kubernetes/utils.py index 1885bce04..518b0e96d 100644 --- a/kubeflow/spark/backends/kubernetes/utils.py +++ b/kubeflow/spark/backends/kubernetes/utils.py @@ -35,6 +35,7 @@ from kubeflow.spark.types.types import ( Driver, Executor, + ResourceDict, SparkConnectInfo, SparkConnectState, SparkJob, @@ -116,6 +117,8 @@ def _resolve_driver_resources( Raises: ValueError: If the configured CPU or memory values are invalid. + TypeError: + If the configured CPU value is not a supported type. """ cores = constants.DEFAULT_DRIVER_CPU @@ -127,7 +130,7 @@ def _resolve_driver_resources( if "memory" in driver.resources: memory = _memory_kubernetes_to_spark( - driver.resources["memory"], + str(driver.resources["memory"]), ) return cores, memory @@ -136,7 +139,7 @@ def _resolve_driver_resources( def _resolve_executor_resources( executor: Executor | None = None, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, ) -> tuple[int, int, str]: """Resolve executor configuration. @@ -148,7 +151,7 @@ def _resolve_executor_resources( Number of executor instances. resources_per_executor: - Resource requirements. + Resource requirements as ResourceDict. Returns: Tuple containing ``(instances, cores, memory)``. @@ -156,6 +159,8 @@ def _resolve_executor_resources( Raises: ValueError: If the configured CPU or memory values are invalid. + TypeError: + If the configured CPU value is not a supported type. """ if executor and executor.num_instances is not None: @@ -183,7 +188,7 @@ def _resolve_executor_resources( if "memory" in resource_dict: memory = _memory_kubernetes_to_spark( - resource_dict["memory"], + str(resource_dict["memory"]), ) return instances, cores, memory @@ -247,22 +252,26 @@ def _memory_kubernetes_to_spark(memory: str) -> str: return f"{math.ceil(total_bytes / (2**20))}m" -def _validate_cpu_value(cpu: str | int | None) -> int: +def _validate_cpu_value(cpu: str | int | float | None) -> int: """Validate and normalize CPU cores value. Args: - cpu: CPU value provided by user. + cpu: CPU value provided by user (str, int, or float). Returns: Integer CPU core value. Raises: - ValueError: If CPU value is invalid. + ValueError: If CPU value is invalid or non-positive. + TypeError: If CPU value is not a supported type or is a boolean. """ if cpu is None: raise ValueError("CPU value cannot be None") - if isinstance(cpu, int): + if isinstance(cpu, bool): + raise TypeError(f"Invalid CPU type '{type(cpu).__name__}'. Expected str, int, or float.") + + if isinstance(cpu, (int, float)): cores = float(cpu) elif isinstance(cpu, str): @@ -282,10 +291,13 @@ def _validate_cpu_value(cpu: str | int | None) -> int: cores = int(milli_cpu) / 1000 else: - cores = float(cpu) + try: + cores = float(cpu) + except ValueError as e: + raise ValueError(f"Invalid CPU value: {cpu!r}") from e else: - raise ValueError(f"Invalid CPU type '{type(cpu)}'. Expected str or int.") + raise TypeError(f"Invalid CPU type '{type(cpu).__name__}'. Expected str, int, or float.") if not math.isfinite(cores) or cores <= 0: raise ValueError(f"Invalid CPU value: {cpu!r}") @@ -315,12 +327,12 @@ def apply_options( backend: Backend used for option validation. - Raises: - ValueError: - If options are provided without a backend instance. + Raises: + ValueError: + If options are provided without a backend instance. - TypeError: - If an option is not callable. + TypeError: + If an option is not callable. """ if not options: return @@ -450,6 +462,8 @@ def get_spark_connect_driver_spec( Raises: ValueError: If the configured driver resources are invalid. + TypeError: + If the configured driver CPU resource is not a supported type. """ cores, memory = _resolve_driver_resources(driver) @@ -473,7 +487,7 @@ def get_spark_connect_driver_spec( def get_spark_connect_executor_spec( executor: Executor | None = None, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, ) -> models.SparkV1alpha1ExecutorSpec: """Convert SDK Executor to API ExecutorSpec. @@ -484,7 +498,7 @@ def get_spark_connect_executor_spec( Args: executor: SDK Executor configuration. num_executors: Simple mode number of executors. - resources_per_executor: Simple mode resource requirements. + resources_per_executor: Simple mode resource requirements as ResourceDict. Returns: API ExecutorSpec model. @@ -492,6 +506,8 @@ def get_spark_connect_executor_spec( Raises: ValueError: If the configured executor resources are invalid. + TypeError: + If the configured executor CPU resource is not a supported type. """ instances, cores, memory = _resolve_executor_resources( executor, @@ -511,7 +527,7 @@ def build_spark_connect_cr( namespace: str, spark_version: str | None = None, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, spark_conf: dict[str, str] | None = None, driver: Driver | None = None, executor: Executor | None = None, @@ -531,7 +547,7 @@ def build_spark_connect_cr( namespace: Kubernetes namespace. spark_version: Spark version (default: `constants.DEFAULT_SPARK_VERSION`). num_executors: Number of executor instances (simple mode). - resources_per_executor: Resource requirements per executor (simple mode). + resources_per_executor: Resource requirements per executor as ResourceDict (simple mode). spark_conf: Spark configuration properties. driver: Driver configuration (advanced mode). executor: Executor configuration (advanced mode). @@ -544,6 +560,8 @@ def build_spark_connect_cr( Raises: ValueError: If the provided driver or executor resource configuration is invalid. + TypeError: + If the provided driver or executor CPU resource is not a supported type. """ _validate_spark_conf(spark_conf) @@ -676,6 +694,8 @@ def get_spark_job_driver_spec( Raises: ValueError: If the default driver resource configuration is invalid. + TypeError: + If the driver CPU resource is not a supported type. """ cores, memory = _resolve_driver_resources(driver) @@ -688,13 +708,13 @@ def get_spark_job_driver_spec( def get_spark_job_executor_spec( num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, ) -> models.SparkV1beta2ExecutorSpec: """Build ExecutorSpec for SparkApplication. Args: num_executors: Number of executor instances. - resources_per_executor: Resource requirements for each executor. + resources_per_executor: Resource requirements for each executor as ResourceDict. Returns: SparkApplication ExecutorSpec model. @@ -702,6 +722,8 @@ def get_spark_job_executor_spec( Raises: ValueError: If the configured executor resources are invalid. + TypeError: + If the executor CPU resource is not a supported type. """ instances, cores, memory = _resolve_executor_resources( num_executors=num_executors, @@ -802,7 +824,7 @@ def get_spark_application_cr_from_file_job( main_file: str, arguments: list[str] | None = None, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, options: list | None = None, backend: Any | None = None, spark_conf: dict[str, str] | None = None, @@ -815,7 +837,7 @@ def get_spark_application_cr_from_file_job( main_file: Path or URI to the Spark application file. arguments: Command-line arguments passed to the Spark application. num_executors: Number of executor instances. - resources_per_executor: Resource requirements for each executor. + resources_per_executor: Resource requirements for each executor as ResourceDict. options: List of configuration options. backend: Backend instance used for option validation. spark_conf: Spark configuration properties. @@ -826,6 +848,8 @@ def get_spark_application_cr_from_file_job( Raises: ValueError: If the executor resource configuration is invalid. + TypeError: + If the executor CPU resource is not a supported type. """ spark_application = models.SparkV1beta2SparkApplication( api_version=f"{constants.SPARK_APPLICATION_GROUP}/{constants.SPARK_APPLICATION_VERSION}", @@ -865,7 +889,7 @@ def get_spark_application_cr_from_func_job( func: Callable, func_args: dict[str, Any] | None = None, num_executors: int | None = None, - resources_per_executor: dict[str, str] | None = None, + resources_per_executor: ResourceDict | None = None, options: list | None = None, backend: Any | None = None, spark_conf: dict[str, str] | None = None, @@ -878,7 +902,7 @@ def get_spark_application_cr_from_func_job( func: Python function to execute as the Spark application. func_args: Keyword arguments passed to the function. num_executors: Number of executor instances. - resources_per_executor: Resource requirements for each executor. + resources_per_executor: Resource requirements for each executor as ResourceDict. options: List of configuration options. backend: Backend instance used for option validation. spark_conf: Spark configuration properties. @@ -890,6 +914,8 @@ def get_spark_application_cr_from_func_job( ValueError: If the provided function is invalid or the executor resource configuration is invalid. + TypeError: + If the executor CPU resource is not a supported type. """ _validate_spark_conf(spark_conf) diff --git a/kubeflow/spark/types/types.py b/kubeflow/spark/types/types.py index e1779139d..9be9829bd 100644 --- a/kubeflow/spark/types/types.py +++ b/kubeflow/spark/types/types.py @@ -21,6 +21,9 @@ import logging from typing import Any +# Type alias for driver and executor resource dictionary specifications +ResourceDict = dict[str, str | int | float] + logger = logging.getLogger(__name__) @@ -87,7 +90,7 @@ class Driver: """ image: str | None = None - resources: dict[str, str] | None = None + resources: ResourceDict | None = None java_options: str | None = None service_account: str | None = None @@ -121,7 +124,7 @@ class Executor: """ num_instances: int | None = None - resources_per_executor: dict[str, str] | None = None + resources_per_executor: ResourceDict | None = None java_options: str | None = None