From c65030de927369685ba2b00135fd8aa5a7bf2d66 Mon Sep 17 00:00:00 2001 From: Amogh Desai Date: Mon, 31 Aug 2026 06:16:30 -0500 Subject: [PATCH] Restructure spark durable execution docs for all cluster managers --- providers/apache/spark/docs/operators.rst | 90 ++++++++++++++--------- 1 file changed, 56 insertions(+), 34 deletions(-) diff --git a/providers/apache/spark/docs/operators.rst b/providers/apache/spark/docs/operators.rst index 6da0e87fce66a..92792a2672526 100644 --- a/providers/apache/spark/docs/operators.rst +++ b/providers/apache/spark/docs/operators.rst @@ -182,37 +182,41 @@ Reference For further information, look at `Apache Spark submitting applications `_. -Cluster mode crash recovery (Spark standalone) -""""""""""""""""""""""""""""""""""""""""""""""" +Durable execution (crash recovery) +"""""""""""""""""""""""""""""""""" -When running in Spark standalone cluster mode (``--deploy-mode cluster``), the Spark driver runs -independently on the cluster. If the Airflow worker dies while the Spark job is running, the driver keeps running but -Airflow loses track of it and the behaviour to submit a brand new job would be wasting -the compute already done or even cause conflicts if the Spark job itself is not designed to be idempotent. +In cluster deploy mode (``--deploy-mode cluster``) the Spark driver runs independently of the +Airflow worker. If the worker dies while the Spark job is running, the driver keeps running but +Airflow loses track of it, and submitting a brand new job on retry wastes the compute already done, +or even causes conflicts if the Spark job itself is not designed to be idempotent. -Now, the ``SparkSubmitOperator`` solves this by persisting the driver ID to :doc:`task state store +The ``SparkSubmitOperator`` solves this by persisting the driver identifier to :doc:`task state store ` immediately after submission. On retry, it reads -the ID back and reconnects to the already-running driver instead of resubmitting. +the identifier back and reconnects to the already-running driver instead of resubmitting. This is the **synchronous path** — the worker holds a slot for the duration of polling. This is a crash-safety net for teams running sync operators for log observability, org constraints, or because a Triggerer is not available. Teams with a Triggerer available may also consider deferrable operators, which free the worker slot but may come with added complexity. -**Connection requirements for crash recovery** +``durable`` defaults to ``True`` and every cluster manager supports it. What differs is the +prerequisite each one needs before the operator can check driver status on retry: -The reconnection polling calls the Spark standalone REST API -(``GET /v1/submissions/status/{driverId}``). Make sure the Spark connection's -``REST scheme`` and ``REST port`` extras match your cluster's configuration: - -* ``REST scheme`` — set to ``https`` if your cluster has TLS enabled on the REST port - (``spark.ssl.standalone.enabled=true``). Defaults to ``http``. -* ``REST port`` — set to the value of ``spark.master.rest.port`` on your cluster. Defaults to ``6066``. +.. list-table:: + :header-rows: 1 + :widths: 30 70 -See :doc:`connections/spark-submit` for how to configure these fields. + * - Cluster manager + - Prerequisite for ``durable=True`` + * - `Spark standalone`_ + - ``REST scheme`` and ``REST port`` connection extras + * - `Kubernetes cluster mode`_ + - ``track_driver_via_k8s_api=True`` + * - `YARN cluster mode`_ + - ``yarn_track_via_rm_api=True`` and ``yarn_resourcemanager_webapp_address`` .. note:: - Crash recovery in cluster mode requires Airflow 3.3+ (``task_state_store`` support). Below + Durable execution requires Airflow 3.3+ (``task_state_store`` support). Below 3.3, ``durable`` has no effect: setting it explicitly only emits a warning, and the operator always submits fresh, exactly as before this feature existed. The deprecated ``reconnect_on_retry`` parameter (the original name for this same feature, superseded almost @@ -229,8 +233,25 @@ This is most reliable for deferred tasks (``deferrable=True``); clearing a task polling synchronously can cancel the driver via ``on_kill`` before the next attempt gets a chance to reconnect -- see :doc:`apache-airflow:core-concepts/resumable-tasks` for why. -Tracking driver status via Kubernetes API -"""""""""""""""""""""""""""""""""""""""""" +.. _cluster-mode-crash-recovery-spark-standalone: + +Spark standalone +"""""""""""""""" + +Reconnecting to a running driver calls the Spark standalone REST API +(``GET /v1/submissions/status/{driverId}``). Make sure the Spark connection's +``REST scheme`` and ``REST port`` extras match your cluster's configuration: + +* ``REST scheme`` — set to ``https`` if your cluster has TLS enabled on the REST port + (``spark.ssl.standalone.enabled=true``). Defaults to ``http``. +* ``REST port`` — set to the value of ``spark.master.rest.port`` on your cluster. Defaults to ``6066``. + +See :doc:`connections/spark-submit` for how to configure these fields. + +.. _tracking-driver-status-via-kubernetes-api: + +Kubernetes cluster mode +""""""""""""""""""""""" When running in Kubernetes cluster mode, ``spark-submit`` blocks for the duration of the job. The JVM runs processes which does nothing but polling of the pod phase and holds heap space for @@ -238,7 +259,10 @@ the entire duration. This is not ideal for long-running jobs, especially when th for long periods (e.g. waiting for data or user input). Set ``track_driver_via_k8s_api=True`` to have the operator track the driver pod status via the -Python Kubernetes client rather than holding ``spark-submit`` open for the full job duration: +Python Kubernetes client rather than holding ``spark-submit`` open for the full job duration. The +same flag is what lets the operator find the driver again after a crash, so it is also the +prerequisite for durable execution on Kubernetes: the driver pod name is persisted to task state +before polling begins, and a retry reconnects to that pod instead of submitting a fresh one. .. code-block:: python @@ -260,17 +284,17 @@ Python Kubernetes client rather than holding ``spark-submit`` open for the full conflicts with the flag and a ``ValueError`` will be raised at task start. * The Airflow worker must be able to reach the Kubernetes API server and have permission to read and delete pods in the driver's namespace; otherwise pod tracking and cleanup will fail. -* Set ``durable=True`` (the default) to enable crash recovery: the driver pod name is - persisted to task state before polling begins, so a worker crash and retry reconnects to the - existing pod instead of submitting a fresh one. Set ``durable=False`` to always - submit a fresh driver on retry. * Pod completion is detected from ``pod.status.phase``. If your driver pods have sidecar containers (e.g. Istio injection enabled for the driver namespace), the pod phase may not advance to ``Succeeded`` until all sidecars exit. In that case the poll loop will wait indefinitely — set ``execution_timeout`` as a hard bound. -YARN ResourceManager API tracking -""""""""""""""""""""""""""""""""" +Set ``durable=False`` to always submit a fresh driver on retry. + +.. _yarn-resourcemanager-api-tracking: + +YARN cluster mode +""""""""""""""""" When running Spark applications on YARN in cluster deploy mode, the default Spark submit path keeps the local ``spark-submit`` JVM alive on the Airflow worker while the YARN @@ -282,6 +306,11 @@ application, then poll the YARN ResourceManager REST API until the application r state. The ResourceManager API polling interval is controlled by ``status_poll_interval`` with a minimum of 10 seconds. +The ResourceManager REST API is also what makes checking application status on retry possible, so +this flag is the prerequisite for durable execution on YARN. Because ``durable`` defaults to +``True``, leaving the flag unset raises a ``ValueError`` at task start rather than silently falling +back to a fire-and-forget submission. + This mode requires the Spark connection extra to set ``yarn_resourcemanager_webapp_address`` before the application is submitted: @@ -306,13 +335,6 @@ the application is submitted: yarn_track_via_rm_api=True, ) -On Airflow 3.3+, YARN cluster mode with ``durable=True`` (the default) requires -``yarn_track_via_rm_api=True`` -- the ResourceManager REST API is what makes checking application -status on retry possible. Without it, the operator raises a ``ValueError`` at task start rather -than silently falling back to a fire-and-forget submission. Below 3.3, ``durable`` has no effect -at all, so this requirement doesn't apply there either: durable execution isn't active to have a -prerequisite for. - For Kerberized clusters, install ``requests-kerberos`` in the Airflow environment. When the Spark connection has both ``keytab`` and ``principal`` configured, Airflow automatically uses ``HTTPKerberosAuth()`` for the ResourceManager REST requests.