simple test DAG
This commit is contained in:
parent
0a62276c42
commit
ba99672349
|
@ -64,7 +64,7 @@ dag = DAG(
|
|||
|
||||
submit = SparkKubernetesOperator(
|
||||
task_id='spark_pi_submit',
|
||||
namespace='spark-jobs',
|
||||
namespace='lot1-spark-jobs',
|
||||
application_file="example_spark_kubernetes_operator_pi.yaml",
|
||||
kubernetes_conn_id="kubernetes_default",
|
||||
do_xcom_push=True,
|
||||
|
@ -73,7 +73,7 @@ submit = SparkKubernetesOperator(
|
|||
|
||||
sensor = SparkKubernetesSensor(
|
||||
task_id='spark_pi_monitor',
|
||||
namespace='spark-jobs',
|
||||
namespace='lot1-spark-jobs',
|
||||
application_name="{{ task_instance.xcom_pull(task_ids='spark_pi_submit')['metadata']['name'] }}",
|
||||
kubernetes_conn_id="kubernetes_default",
|
||||
dag=dag,
|
||||
|
|
Loading…
Reference in New Issue