与您对问题的评论类似,我在回填大型数据库时解决此问题的方法是让 dag 生成器基于 connection_created_on 和 @ 创建三个 dag(两个回填和一个正在进行) 987654322@ 值。
正在进行的 dag 每小时运行一次,并在 connection_created_on 值的同一天午夜开始。这两个回填然后从当月的第一天开始每天拉动,然后从start_date 的第一个月开始每月拉动。在这种情况下,我知道我们总是希望从每月的第一天开始,并且最多一个月的数据足够小,可以汇总在一起,因此为了方便起见,我将其分为这三种 dag 类型。
def create_dag(dag_id,
schedule,
db_conn_id,
default_args,
catchup=False,
max_active_runs=3):
dag = DAG(dag_id,
default_args=default_args,
schedule_interval=schedule,
catchup=catchup,
max_active_runs=max_active_runs
)
with dag:
kick_off_dag = DummyOperator(task_id='kick_off_dag')
return dag
db_conn_id = 'my_first_db_conn'
connection_created_on = '2018-05-17 12:30:54.271Z'
hourly_id = '{}_to_redshift_hourly'.format(db_conn_id)
daily_id = '{}_to_redshift_daily_backfill'.format(db_conn_id)
monthly_id = '{}_to_redshift_monthly_backfill'.format(db_conn_id)
start_date = '2005-01-01 00:00:00.000Z'
start_date = datetime.strptime(start_date, '%Y-%m-%dT%H:%M:%S.%fZ')
start_date = datetime(start_date.year, start_date.month, 1)
cco_datetime = datetime.strptime(connection_created_on, '%Y-%m-%dT%H:%M:%S.%fZ')
hourly_start_date = datetime(cco_datetime.year, cco_datetime.month, cco_datetime.day)
daily_start_date = hourly_start_date - timedelta(days=(cco_datetime.day-1))
daily_end_date = hourly_start_date - timedelta(days=1)
monthly_start_date = start_date if start_date else hourly_start_date - timedelta(days=365+cco_datetime.day-1)
monthly_end_date = daily_start_date
globals()[hourly_id] = create_dag(hourly_id,
'@hourly',
db_conn_id,
{'start_date': hourly_start_date,
'retries': 2,
'retry_delay': timedelta(minutes=5),
'email': [],
'email_on_failure': True,
'email_on_retry': False},
catchup=True,
max_active_runs=1)
globals()[daily_id] = create_dag(daily_id,
'@daily',
db_conn_id,
{'start_date': daily_start_date,
'end_date': daily_end_date,
'retries': 2,
'retry_delay': timedelta(minutes=5),
'email': [],
'email_on_failure': True,
'email_on_retry': False},
catchup=True)
globals()[monthly_id] = create_dag(monthly_id,
'@monthly',
db_conn_id,
{'start_date': monthly_start_date,
'end_date': monthly_end_date,
'retries': 2,
'retry_delay': timedelta(minutes=5),
'email': [],
'email_on_failure': True,
'email_on_retry': False},
catchup=True)