Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
90 changes: 56 additions & 34 deletions providers/apache/spark/docs/operators.rst
Original file line number Diff line number Diff line change
Expand Up @@ -182,37 +182,41 @@ Reference

For further information, look at `Apache Spark submitting applications <https://spark.apache.org/docs/latest/submitting-applications.html>`_.

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
<apache-airflow:core-concepts/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
Expand All @@ -229,16 +233,36 @@ 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
the entire duration. This is not ideal for long-running jobs, especially when the driver is idle
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

Expand All @@ -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
Expand All @@ -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:

Expand All @@ -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.
Expand Down