【问题标题】:Airflow keeps running my DAG, despite catchup=False, schedule_interval=datetime.timedelta(hours=2)尽管 catchup=False, schedule_interval=datetime.timedelta(hours=2),Airflow 仍在运行我的 DAG
【发布时间】:2019-07-30 23:26:38
【问题描述】:

与之前的问题类似,但给出的答案均无效。我有一个 DAG:

import datetime
import os

from airflow import DAG
from airflow.contrib.operators.dataflow_operator import DataflowTemplateOperator
from airflow.operators import BashOperator

PROJECT = os.environ['PROJECT']
GCS_BUCKET = os.environ['BUCKET']
API_KEY = os.environ['API_KEY']

default_args = {
    'owner': 'me',
    'start_date': datetime.datetime(2019, 7, 30),
    'depends_on_past': False,
    'email': [''],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 0,
    'retry_delay': datetime.timedelta(hours=1),
    'catchup': False
}

dag = DAG('dag-name',
          schedule_interval=datetime.timedelta(hours=2),
          default_args=default_args,
          max_active_runs=1,
          concurrency=1,
          catchup=False)

DEFAULT_OPTIONS_TEMPLATE = {
    'project': PROJECT,
    'stagingLocation': 'gs://{}/staging'.format(GCS_BUCKET),
    'tempLocation': 'gs://{}/temp'.format(GCS_BUCKET)
}

def my-dataflow-job(template_location, name):
    run_time = datetime.datetime.utcnow()
    a_value = run_time.strftime('%Y%m%d%H')

    t1 = DataflowTemplateOperator(
        task_id='{}-task'.format(name),
        template=template_location,
        parameters={'an_argument': a_value},
        dataflow_default_options=DEFAULT_OPTIONS_TEMPLATE,
        poll_sleep=30
    )

    t2 = BashOperator(
        task_id='{}-loader-heartbeat'.format(name),
        bash_command='curl --fail -XGET "[a heartbeat URL]" --header "Authorization: heartbeat_service {1}"'.format(name, API_KEY)
    )

    t1 >> t2

with dag:
    backup_bt_to_bq('gs://[path to gcs]'.format(GCS_BUCKET), 'name')

如您所见,我正在努力阻止 Airflow 尝试回填。然而,当我部署 DAG 时(当天晚些时候,2019 年 7 月 30 日),它只是一个接一个、一个接一个、一个接一个地运行 DAG。

由于此任务正在移动一些数据,因此这是不可取的。如何让气流遵守“每隔一小时运行一次” schedule_interval??

如您所见,我在 DAG args 和默认 args 中都设置了 catchup: False(以防万一,在 DAG args 中从它们开始)。重试延迟也很长。

每次 DAG 运行都会报告为成功。 我正在使用以下版本:

composer-1.5.0-airflow-1.10.1

我的下一步是 Kubernetes cron...

【问题讨论】:

    标签: airflow-scheduler


    【解决方案1】:

    我怀疑你第一次创建 dag 时没有catchup=False。我认为气流可能无法识别初始 dag 创建后追赶参数的变化。

    尝试重命名它,看看会发生什么。例如。添加 v2 并启用它。启用后,即使 catchup 为 false,它也会运行一次,因为有一个有效的完成时间间隔(即当前时间 >= start_time + schedule_interval),但仅此而已。

    当然,使用不会做任何昂贵操作的假操作员进行测试。

    【讨论】:

    • 我确实启用了它。并尝试了几个名字。可能只是转储到 k8s cron 中。
    • 也许如果您更新问题中的代码,我可能会发现问题。很明显,您发布的代码不是最新的:my-dataflow-job 不是有效的函数名,backup_bt_to_bq 未在任何地方定义。
    猜你喜欢
    • 2018-05-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多