# coding=utf-8 """ @Time :2026/6/24 @Author :tao.chen """ from datetime import datetime from pydantic import BaseModel, Field, field_validator class Job(BaseModel): job_id: str application_id: str script_path: str queue: str submit_time: datetime connection: str yarn_rm_url: str | None = None class JobStatus(BaseModel): application_id: str state: str raw: str = Field(default="") class JobResult(BaseModel): application_id: str state: str final_status: str | None = None diagnostics: str | None = None tracking_url: str | None = None started_time: int | None = None finished_time: int | None = None class SubmitResult(BaseModel): job_id: str application_id: str tracking_url: str | None = None class Connection(BaseModel): name: str # Defaults to "yarn" because that's the literal string spark-submit wants # for --master when targeting YARN. Override for Standalone (spark://...), # Kubernetes (k8s://...), or local mode. master: str = "yarn" deploy_mode: str = "cluster" yarn_rm_url: str | None = None spark_conf: dict[str, str] = Field(default_factory=dict) # None means "fall back to Settings.ssl_verify_default". Explicit True/False # overrides the global default for this connection. ssl_verify: bool | None = None ssl_ca_bundle: str | None = None # Authentication for YARN REST calls. auth_type: str = "none" # "none" | "simple" | "basic" | "kerberos" auth_user: str | None = None auth_password: str | None = None # Display/audit only for kerberos; actual SPNEGO uses the system cache. auth_principal: str | None = None auth_keytab: str | None = None @field_validator("master") @classmethod def _check_master(cls, v: str) -> str: """Catch common typos like 'yarn-cluster' or 'http://...'. """ if v == "yarn": return v if v.startswith(("spark://", "k8s://", "mesos://", "local")): return v raise ValueError( f"master must be 'yarn', 'spark://...', 'k8s://...', 'mesos://...', " f"or 'local[/N]'; got {v!r}" ) @field_validator("auth_type") @classmethod def _check_auth_type(cls, v: str) -> str: if v not in {"none", "simple", "basic", "kerberos"}: raise ValueError( f"auth_type must be one of none/simple/basic/kerberos; got {v!r}" ) return v class PendingSubmission(BaseModel): pending_id: str app_name: str | None = None connection: str master: str deploy_mode: str yarn_rm_url: str | None = None script_path: str queue: str executor_memory: str executor_cores: int num_executors: int spark_conf: dict[str, str] = Field(default_factory=dict) extra_args: dict[str, str] = Field(default_factory=dict) created_at: datetime status: str = "PENDING" # PENDING | SUBMITTED | CANCELLED | FAILED error: str | None = None job_id: str | None = None application_id: str | None = None tracking_url: str | None = None