0
votes

I am trying to trigger DAG task for 3 times , how can this be done using python script.

Currently my flow job is

dag = DAG('dag1', default_args=default_args, concurrency=1, max_active_runs=1, schedule_interval=None, catchup=False)

task = BashOperator( task_id='task1', bash_command='ssh "runJob.sh"', dag=dag)

Is there any better way to trigger job for a given specified number of times ?

1

1 Answers

0
votes

I assume you use Airflow 2, you can use the following command to trigger a DAG run.

airflow dags trigger [-h] [-c CONF] [-e EXEC_DATE] [-r RUN_ID] [-S SUBDIR] dag_id

You can read more on this command here.

Based on that command, we can use use subprocess to run it.

import subprocess
command = [
    "airflow",
    "dags",
    "trigger",
    "-e",
    "2021-02-05 00:00:00",
    "-r",
    "manual__2021-02-05T00:00:00+00:00",
    "my_dag",
]

for i in range(3):
    process = subprocess.Popen(command, stdout=subprocess.PIPE)
    output, error = process.communicate()