I do submit spark application to kubernetes from Airflow as SparkKubernetesOperator task.
When I mark task as failed I see the pod is not deleted in kubernetes.
How to fix it?
dag = DAG(
"dag-name",
default_args=default_args,
description='submit dag',
schedule_interval="30 1 * * *",
start_date = datetime(2022, 7, 26),
)
t1 = SparkKubernetesOperator(
task_id='sample-name',
namespace="batch",
application_file="k8s/sample.yaml",
do_xcom_push=True,
dag=dag,
params = {"processDate": process_date},
)
t1_sensor = SparkKubernetesSensor(
task_id='sample-monitor',
namespace="batch",
application_name="{{ task_instance.xcom_pull(task_ids='sample-name')['metadata']['name'] }}",
kubernetes_conn_id="kubernetes_default",
dag=dag,
)
t1 >> t1_sensor