from datetime import datetime, timedelta from airflow import DAG from airflow.operators.bash_operator import BashOperator from airflow.operators.dummy_operator import DummyOperator jar_file_path = "/path/to/yourJarFile.jar" input_table_name = "your_input_hive_table_name" output_table_name = "your_output_hive_table_name" default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 1, 1), 'depends_on_past': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } dag = DAG( 'my_data_pipeline_dag', default_args=default_args, description='DAG to run the MyExampleDataPipeline class', schedule_interval='@daily', # Adjust the schedule as needed ) start_task = DummyOperator(task_id='start', dag=dag) end_task = DummyOperator(task_id='end', dag=dag) spark_submit_command = f"spark-submit --class MyExampleDataPipeline --master local[*] {jar_file_path} {input_table_name} {output_table_name}" run_data_pipeline_task = BashOperator( task_id='run_data_pipeline', bash_command=spark_submit_command, dag=dag, ) start_task >> run_data_pipeline_task >> end_task