[ 
https://issues.apache.org/jira/browse/SPARK-58119?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-58119:
-----------------------------------
    Labels: pull-request-available  (was: )

> [CORE][PYTHON] Introduce a PythonWorkerHandle abstraction for the Python 
> worker path
> ------------------------------------------------------------------------------------
>
>                 Key: SPARK-58119
>                 URL: https://issues.apache.org/jira/browse/SPARK-58119
>             Project: Spark
>          Issue Type: Improvement
>          Components: PySpark
>    Affects Versions: 4.1.2
>            Reporter: Fabian Paul
>            Priority: Major
>              Labels: pull-request-available
>
> h3. Problem
> The Python worker execution path (\{{SparkEnv.createPythonWorker}}, 
> \{{PythonWorkerFactory}}, \{{BasePythonRunner}}/\{{PythonRunner}}, 
> \{{PythonUDFRunner}}, \{{PythonArrowOutput}}, \{{PythonPlannerRunner}}) uses 
> \{{java.lang.ProcessHandle}} directly to represent a launched worker. This 
> couples the runner to a sealed JDK type, and the small amount of behavior the 
> runner actually needs from a worker handle (liveness, termination, 
> post-mortem diagnostics) is not expressed as a first-class Spark type. In 
> particular, the Python faulthandler-log reading is duplicated inline at each 
> crash / idle-timeout site rather than living behind a single, testable method.
> h3. Goal
> Introduce a Spark-owned \{{PythonWorkerHandle}} trait that exposes only the 
> operations the runner needs:
> * isAlive(): Boolean
> * destroy(): Boolean
> * terminationDiagnostics(faultHandlerEnabled: Boolean): Option[String]
> The default backend, \{{LocalPythonWorkerHandle}}, wraps a 
> \{{java.lang.ProcessHandle}} and delegates 1:1; a 
> \{{PythonWorkerHandle.of(pid, faultHandlerLog)}} factory constructs it. The 
> runner supplies the (SQLConf-resolved) faulthandler flag when it reads 
> diagnostics, and the handle locates its own worker's faulthandler log, so the 
> reader no longer needs to know about faulthandler log files. The OS pid is 
> deliberately kept out of the handle contract: its one consumer, the daemon 
> kill-by-pid protocol, is served from a factory-internal map.
> h3. Impact
> Internal refactor only. No user-facing behavior, configuration, or API 
> change: \{{LocalPythonWorkerHandle}} delegates to the JDK handle and the 
> existing faulthandler read, and the diagnostics path is behavior-preserving. 
> Decouples the runner from the concrete JDK type and centralizes the "what the 
> runner needs from a worker handle" contract, including the diagnostics read, 
> in one place.
> h3. Approach
> * Add \{{PythonWorkerHandle}} (trait), \{{LocalPythonWorkerHandle}} (impl), 
> and \{{PythonWorkerHandle.of}} in \{{core/.../api/python/}}.
> * Change \{{SparkEnv.createPythonWorker}} and the reader plumbing from 
> \{{Option[ProcessHandle]}} to \{{Option[PythonWorkerHandle]}}.
> * Route the crash / idle-timeout diagnostics through 
> \{{handle.terminationDiagnostics(faultHandlerEnabled)}}.
> * Add \{{PythonWorkerHandleSuite}}.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to