# coding=utf-8 """ @Time :2026/6/24 @Author :tao.chen """ import secrets import uuid from datetime import datetime from common.logging import logger from spark_executor.core.connection_store import store as conn_store from spark_executor.core.job_store import JobStore from spark_executor.core.log_parser import parse_spark_submit_output from spark_executor.core.pending_store import store as pending_store from spark_executor.core.spark_submit import ( SparkSubmitError, build_spark_submit_command, run_spark_submit, ) from spark_executor.models import Job, PendingSubmission, SubmitResult def _new_pending_id() -> str: return "p_" + secrets.token_hex(6) def prepare_submit_job( *, connection: str, script_path: str, queue: str = "default", executor_memory: str = "4G", executor_cores: int = 2, num_executors: int = 2, ) -> dict[str, object]: """Snapshot connection params and persist a PendingSubmission. Does NOT submit.""" logger.debug( f"prepare_submit_job enter connection={connection} script_path={script_path} " f"queue={queue} executor_memory={executor_memory} executor_cores={executor_cores} " f"num_executors={num_executors}" ) conn = conn_store.get(connection) if conn is None: raise KeyError(f"Unknown connection: {connection}") pending_id = _new_pending_id() pending = PendingSubmission( pending_id=pending_id, connection=connection, master=conn.master, deploy_mode=conn.deploy_mode, script_path=script_path, queue=queue, executor_memory=executor_memory, executor_cores=executor_cores, num_executors=num_executors, spark_conf=dict(conn.spark_conf), created_at=datetime.utcnow(), status="PENDING", ) pending_store.save(pending) logger.info( f"prepare_submit_job ok pending_id={pending_id} connection={connection} " f"master={conn.master} script_path={script_path}" ) return { "pending_id": pending_id, "status": "PENDING", "parameters": pending.model_dump(), } # Module-level job store singleton; replaced in tests. job_store: JobStore = JobStore() def confirm_submit_job(*, pending_id: str) -> SubmitResult: """Actually invoke spark-submit for a previously-prepared PendingSubmission.""" logger.debug(f"confirm_submit_job enter pending_id={pending_id}") pending = pending_store.get(pending_id) if pending is None: raise KeyError(f"Unknown pending_id: {pending_id}") if pending.status != "PENDING": raise ValueError( f"pending_id {pending_id} is in status {pending.status!r}, not PENDING" ) cmd = build_spark_submit_command( master=pending.master, deploy_mode=pending.deploy_mode, script_path=pending.script_path, queue=pending.queue, executor_memory=pending.executor_memory, executor_cores=pending.executor_cores, num_executors=pending.num_executors, spark_conf=pending.spark_conf, ) logger.info( f"confirm_submit_job start pending_id={pending_id} " f"application_target={pending.master} script_path={pending.script_path}" ) try: result = run_spark_submit(cmd) except SparkSubmitError as exc: pending.status = "FAILED" pending.error = str(exc) pending_store.save(pending) logger.error(f"confirm_submit_job failed pending_id={pending_id} err={exc}") raise application_id, tracking_url = parse_spark_submit_output(result.stderr) job_id = uuid.uuid4().hex[:12] job_store.put( Job( job_id=job_id, application_id=application_id, script_path=pending.script_path, queue=pending.queue, submit_time=datetime.utcnow(), connection=pending.connection, ) ) pending.status = "SUBMITTED" pending.job_id = job_id pending.application_id = application_id pending_store.save(pending) logger.info( f"confirm_submit_job ok pending_id={pending_id} job_id={job_id} " f"application_id={application_id}" ) return SubmitResult( job_id=job_id, application_id=application_id, tracking_url=tracking_url, ) def list_pending_jobs() -> list[dict[str, object]]: logger.debug("list_pending_jobs enter") return [p.model_dump() for p in pending_store.list_all()] def get_pending_job(pending_id: str) -> dict[str, object]: logger.debug(f"get_pending_job enter pending_id={pending_id}") p = pending_store.get(pending_id) if p is None: raise KeyError(f"Unknown pending_id: {pending_id}") return p.model_dump() def cancel_pending_job(pending_id: str) -> dict[str, str]: logger.debug(f"cancel_pending_job enter pending_id={pending_id}") p = pending_store.get(pending_id) if p is None: raise KeyError(f"Unknown pending_id: {pending_id}") if p.status in ("SUBMITTED", "FAILED"): raise ValueError( f"pending_id {pending_id} is in status {p.status!r} and cannot be cancelled" ) p.status = "CANCELLED" pending_store.save(p) logger.info(f"cancel_pending_job ok pending_id={pending_id}") return {"pending_id": pending_id, "status": "CANCELLED"}