【问题标题】:Airflow backfills and new dag runs气流回填和新的 dag 运行
【发布时间】:2019-01-08 01:51:51
【问题描述】:

我有一个 DAG,它从 2015 年 1 月 1 日到今天每天都有“DAG 运行”。 DAG 中的任务不是“过去依赖的”,这意味着在回填期间它们可以按任何日期顺序执行。

如果我需要回填 DAG 中的任务,我会使用 UI 清除所有任务实例(从今天到过去),然后所有 DAG 运行切换到“正在运行”状态,并且任务从 2015 年 1 月 1 日到今天开始回填.任务很耗时,所以即使由多个线程/worker并行执行,回填也只能在几天内完成。

问题在于,在回填完成之前,调度程序不会添加明天、后天等的新“DAG 运行”,因此我们无法按时计算新天的数据。有什么方法可以在新的一天任务到来时确定它们的优先级,并在新一天的任务完成后继续回填?

附:回填也可以使用“气流回填”CLI 完成,但这种方法有其自身的问题,所以目前我对上述回填技术感兴趣。

【问题讨论】:

  • 我建议将这两个 dag 分开。也许根据配置创建两个 dag。这对你有意义吗?
  • 对于任何来到这里的人,我相信在 DAG 上设置 catchup=True 可以解决此问题。

标签: python scheduled-tasks airflow


【解决方案1】:

与您对问题的评论类似,我在回填大型数据库时解决此问题的方法是让 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)

【讨论】:

    猜你喜欢
    • 2016-12-09
    • 1970-01-01
    • 2022-11-02
    • 2016-12-15
    • 1970-01-01
    • 2021-05-03
    • 2019-11-18
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多