【发布时间】: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...
【问题讨论】: