from datetime import timedelta, datetime import airflow from airflow import DAG from airflow.contrib.operators.ssh_operator import SSHOperator default_args = { 'owner': 'royalties', 'depends_on_past': False, 'email': ['mkhan@theorchard.com'], 'email_on_failure': False, 'email_on_retry': False, 'start_date': datetime.now() - timedelta(minutes=20), 'retries': 0 } dag = DAG(dag_id='file_import', default_args=default_args, schedule_interval=None) t1_bash = """ spark-submit transform.py s3://royalties-spark/stage/dsd244.txt s3://royalties-spark/master/transactions-244-2/ """ t1 = SSHOperator( ssh_conn_id='EMR', task_id='test_ssh_operator', command=t1_bash, dag=dag)