airflow 2.3.3 SparkKubernetesOperator

Viewed 84

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
0 Answers
Related