import logging from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from flows.deezer_daily import config from flows.deezer_daily import tasks logger = logging.getLogger(__file__) default_args = dict(owner="deezer") with DAG( dag_id=config.FEED_NAME, default_args=default_args, schedule_interval='@daily', start_date=datetime(2022, 12, 15), catchup=False, max_active_runs=1, dagrun_timeout=timedelta(minutes=15), ) as dag: licensor = 'theorchard' target_s3_path = config.S3['archive_s3_path'].format( s3_bucket=config.ARCHIVE_S3_BUCKET, spec_version=config.SPEC_VERSION[licensor], datestamp="{{ ds }}", licensor=licensor ) drop_file_name = config.DROP_FILE_NAME[licensor].format( date='{{ ds_nodash }}', ) fetch = PythonOperator( task_id='fetch', # retries=3, # retry_delay=timedelta(minutes=5), python_callable=tasks.fetch, op_kwargs=dict( target_s3_path=target_s3_path, drop_file_name=drop_file_name, ) )