diff --git a/CHANGELOG.md b/CHANGELOG.md index bb22846b..3754173e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -46,6 +46,7 @@ - Make operations infallible where appropriate ([#852], [#860]). - Deprecated airflow `3.2.2` ([#865]). - Bump stackable-operator to 0.119.0 ([#868]). +- Docs: remove Spark submit/monitor classes and refer to demo usage instead ([#870]). ### Fixed @@ -79,6 +80,7 @@ [#862]: https://github.com/stackabletech/airflow-operator/pull/862 [#865]: https://github.com/stackabletech/airflow-operator/pull/865 [#868]: https://github.com/stackabletech/airflow-operator/pull/868 +[#870]: https://github.com/stackabletech/airflow-operator/pull/870 ## [26.7.0] - 2026-07-21 diff --git a/docs/modules/airflow/examples/example_spark_kubernetes_operator.py b/docs/modules/airflow/examples/example_spark_kubernetes_operator.py deleted file mode 100644 index d4f7e575..00000000 --- a/docs/modules/airflow/examples/example_spark_kubernetes_operator.py +++ /dev/null @@ -1,62 +0,0 @@ -# -# Licensed to the Apache Software Foundation (ASF) under one -# or more contributor license agreements. See the NOTICE file -# distributed with this work for additional information -# regarding copyright ownership. The ASF licenses this file -# to you under the Apache License, Version 2.0 (the -# "License"); you may not use this file except in compliance -# with the License. You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, -# software distributed under the License is distributed on an -# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -# KIND, either express or implied. See the License for the -# specific language governing permissions and limitations -# under the License. - -from typing import TYPE_CHECKING, Optional, Sequence -from airflow.models import BaseOperator -from airflow.providers.cncf.kubernetes.hooks.kubernetes import KubernetesHook -import json - -if TYPE_CHECKING: - from airflow.utils.context import Context - - -class SparkKubernetesOperator(BaseOperator): # <1> - template_fields: Sequence[str] = ("application_file", "namespace") - template_ext: Sequence[str] = (".yaml", ".yml", ".json") - ui_color = "#f4a460" - - def __init__( - self, - *, - application_file: str, - namespace: Optional[str] = None, - kubernetes_conn_id: str = "kubernetes_in_cluster", # <2> - api_group: str = "spark.stackable.tech", - api_version: str = "v1alpha1", - **kwargs, - ) -> None: - super().__init__(**kwargs) - self.application_file = application_file - self.namespace = namespace - self.kubernetes_conn_id = kubernetes_conn_id - self.api_group = api_group - self.api_version = api_version - self.plural = "sparkapplications" - - def execute(self, context: "Context"): - hook = KubernetesHook(conn_id=self.kubernetes_conn_id) - self.log.info("Creating SparkApplication...") - self.log.info(json.dumps(self.application_file, indent=4)) - response = hook.create_custom_object( - group=self.api_group, - version=self.api_version, - plural=self.plural, - body=self.application_file, - namespace=self.namespace, - ) - return response diff --git a/docs/modules/airflow/examples/example_spark_kubernetes_sensor.py b/docs/modules/airflow/examples/example_spark_kubernetes_sensor.py deleted file mode 100644 index 3c551d0b..00000000 --- a/docs/modules/airflow/examples/example_spark_kubernetes_sensor.py +++ /dev/null @@ -1,77 +0,0 @@ -# -# Licensed to the Apache Software Foundation (ASF) under one -# or more contributor license agreements. See the NOTICE file -# distributed with this work for additional information -# regarding copyright ownership. The ASF licenses this file -# to you under the Apache License, Version 2.0 (the -# "License"); you may not use this file except in compliance -# with the License. You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, -# software distributed under the License is distributed on an -# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -# KIND, either express or implied. See the License for the -# specific language governing permissions and limitations -# under the License. - -from typing import Optional, Dict -from airflow.exceptions import AirflowException -from airflow.sensors.base import BaseSensorOperator -from airflow.providers.cncf.kubernetes.hooks.kubernetes import KubernetesHook - - -class SparkKubernetesSensor(BaseSensorOperator): # <3> - template_fields = ("application_name", "namespace") - # See https://github.com/stackabletech/spark-k8s-operator/pull/460/files#diff-d737837121132af6b60f50279a78464b05dcfd06c05d1d090f4198a5e962b5f6R371 - # Unknown is set immediately so it must be excluded from the failed states. - FAILURE_STATES = "Failed" - SUCCESS_STATES = "Succeeded" - - def __init__( - self, - *, - application_name: str, - namespace: Optional[str] = None, - kubernetes_conn_id: str = "kubernetes_in_cluster", # <2> - api_group: str = "spark.stackable.tech", - api_version: str = "v1alpha1", - poke_interval: float = 60, - **kwargs, - ) -> None: - super().__init__(**kwargs) - self.application_name = application_name - self.namespace = namespace - self.kubernetes_conn_id = kubernetes_conn_id - self.hook = KubernetesHook(conn_id=self.kubernetes_conn_id) - self.api_group = api_group - self.api_version = api_version - self.poke_interval = poke_interval - - def poke(self, context: Dict) -> bool: - self.log.info("Poking: %s", self.application_name) - response = self.hook.get_custom_object( - group=self.api_group, - version=self.api_version, - plural="sparkapplications", - name=self.application_name, - namespace=self.namespace, - ) - try: - application_state = response["status"]["phase"] - except KeyError: - self.log.debug( - f"SparkApplication status could not be established: {response}" - ) - return False - if application_state in self.FAILURE_STATES: - raise AirflowException( - f"SparkApplication failed with state: {application_state}" - ) - elif application_state in self.SUCCESS_STATES: - self.log.info("SparkApplication ended successfully") - return True - else: - self.log.info("SparkApplication is still in state: %s", application_state) - return False diff --git a/docs/modules/airflow/pages/usage-guide/applying-custom-resources.adoc b/docs/modules/airflow/pages/usage-guide/applying-custom-resources.adoc index 19a24da5..1827e537 100644 --- a/docs/modules/airflow/pages/usage-guide/applying-custom-resources.adoc +++ b/docs/modules/airflow/pages/usage-guide/applying-custom-resources.adoc @@ -56,7 +56,7 @@ dags |_ pyspark_pi.yaml ---- -The Spark job calculates the value of pi using one of the example scripts that comes bundled with Spark: +The Spark job calculates the value of `pi` using one of the example scripts that comes bundled with Spark: [source,yaml] ---- @@ -72,16 +72,7 @@ There are two classes that are used to: The classes `SparkKubernetesOperator` and `SparkKubernetesSensor` are located in two different Python modules as they are typically used for all custom resources and thus are best decoupled from the DAG that calls them. This also demonstrates that modularized DAGs can be used for Airflow jobs as long as all dependencies exist in or below the root folder pulled by git-sync. - -[source,python] ----- -include::example$example_spark_kubernetes_operator.py[] ----- - -[source,python] ----- -include::example$example_spark_kubernetes_sensor.py[] ----- +These files are used in an Airflow demo where they are mounted into the Airflow cluster via a `ConfigMap` https://github.com/stackabletech/demos/blob/main/stacks/airflow/airflow.yaml[here{external-link-icon}^]. [source,python] ---- @@ -97,9 +88,7 @@ include::example$example-spark-dag.py[] [NOTE] ==== -The sensor only tracks the `status.phase` of the SparkApplication. -It cannot copy the Spark driver log into the Airflow task log: the Spark operator deletes the driver Pod as soon as the application reaches `Succeeded` or `Failed`. -To keep driver and executor logs, enable log aggregation for the SparkApplication as described in xref:spark-k8s:usage-guide/logging.adoc[], or use the xref:spark-k8s:usage-guide/history-server.adoc[Spark history server] for event logs. +To persist Spark driver and executor logs independently of any Airflow configuration (such as remote logging), enable log aggregation for the SparkApplication as described in xref:spark-k8s:usage-guide/logging.adoc[], or use the xref:spark-k8s:usage-guide/history-server.adoc[Spark history server] for event logs. ==== Once this DAG is xref:usage-guide/mounting-dags.adoc[mounted] in the DAG folder it can be called and its progress viewed from within the Webserver UI: @@ -110,13 +99,13 @@ Clicking on the "spark_pi_monitor" task and selecting the logs shows that the st image::airflow_dag_log.png[Airflow Connections] -NOTE: If the `KubernetesExecutor` is employed the logs are only accessible via the SDP logging mechanism, described https://docs.stackable.tech/home/stable/concepts/logging[here]. +NOTE: If the `KubernetesExecutor` is employed the logs are only persisted via the SDP logging mechanism, described https://docs.stackable.tech/home/stable/concepts/logging[here] or with xref:usage-guide/logging.adoc[remote logging]. TIP: A full example of the above is used as an integration test https://github.com/stackabletech/airflow-operator/tree/main/tests/templates/kuttl/mount-dags-gitsync[here{external-link-icon}^]. == Logging -As mentioned above, the Airflow task logs are available from the webserver UI if the jobs run with the `celeryExecutor`. +As mentioned above, the Airflow task logs are always available from the webserver UI if the jobs run with the `celeryExecutor` or if xref:usage-guide/logging.adoc[remote logging] has been configured. If the SDP logging mechanism has been deployed, log information can also be retrieved from the vector backend (e.g. Opensearch): image::airflow_dag_log_opensearch.png[Opensearch]