Skip to content
Merged
Show file tree
Hide file tree
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
8 changes: 4 additions & 4 deletions docs/modules/airflow/examples/example-spark-dag.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
from stackable.spark_kubernetes_operator import SparkKubernetesOperator


with DAG( # <4>
with DAG( # <1>
dag_id="sparkapp_dag",
schedule=None,
start_date=datetime(2022, 1, 1),
Expand Down Expand Up @@ -59,20 +59,20 @@ def load_body_to_dict(body):
)
document.update({"metadata": {"name": application_name, "namespace": ns}})

t1 = SparkKubernetesOperator( # <5>
t1 = SparkKubernetesOperator( # <2>
task_id="spark_pi_submit",
namespace=ns,
application_file=document,
do_xcom_push=True,
dag=dag,
)

t2 = SparkKubernetesSensor( # <6>
t2 = SparkKubernetesSensor( # <3>
task_id="spark_pi_monitor",
namespace=ns,
application_name="{{ task_instance.xcom_pull(task_ids='spark_pi_submit')['metadata']['name'] }}",
poke_interval=5,
dag=dag,
)

t1 >> t2 # <7>
t1 >> t2 # <4>
Original file line number Diff line number Diff line change
Expand Up @@ -72,19 +72,18 @@ 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.
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}^].
The two modules are not reproduced here: the repository this layout is taken from is private, as it is used to test git-sync access to protected repositories.
The https://github.com/stackabletech/demos/blob/main/stacks/airflow/airflow.yaml[Airflow demo{external-link-icon}^] shows the same two classes publicly, using a different layout: it defines them in `pyspark_pi.py` itself, before the DAG code, and mounts that file into the Airflow cluster from a `ConfigMap`.
Both classes use the `kubernetes_in_cluster` connection defined above unless a different one is passed as `kubernetes_conn_id`.

[source,python]
----
include::example$example-spark-dag.py[]
----
<1> the wrapper class used for calling the job via `KubernetesHook`
<2> the connection that created for in-cluster usage
<3> the wrapper class used for monitoring the job via `KubernetesHook`
<4> the start of the DAG code
<5> the initial task to invoke the job
<6> the subsequent task to monitor the job
<7> the jobs are chained together in the correct order
<1> the start of the DAG code
<2> the initial task to invoke the job
<3> the subsequent task to monitor the job
<4> the jobs are chained together in the correct order

[NOTE]
====
Expand All @@ -101,8 +100,6 @@ image::airflow_dag_log.png[Airflow Connections]

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 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.
Expand Down
Loading