代码之家  ›  专栏  ›  技术社区  ›  Tomas Jansson

我可以在计划中添加延迟吗?

  •  0
  • Tomas Jansson  · 技术社区  · 7 年前

    我有一个管道,我想每天运行,但我希望执行日期滞后。也就是说,在白天 X 我希望执行日期是 X-3 .有可能吗?

    2 回复  |  直到 7 年前
        1
  •  2
  •   SergiyKolesnikov    7 年前

    看起来你正在使用 execution_date 作为管道逻辑中的变量。例如,要处理比 执行日期 .所以,与其 执行日期 要延迟3天,你可以从 执行日期 并在管道逻辑中使用结果。气流提供了多种方法:

    1. Templates: {{ execution_date - macros.timedelta(days=3) }} .例如 bash_command Bash运算符的参数可以是 bash_command='echo Processing date: {{ execution_date - macros.timedelta(days=3) }} '
    2. The PythonOperator's python callable: 定义可调用函数,比如 def func(execution_date, **kwargs): ... 设置蟒蛇的参数 provide_context=True 这个 执行日期 参数 func() 将设置为当前执行日期( datetime (反对)待命。所以,在里面 func() 你能行 processing_date = execution_date - timedelta(days=3) .
    3. The Sensors' context parameter 当前位置 poke() execute() 任何传感器的方法都有 上下文 包含所有宏的dict参数,包括 执行日期 .所以,在这些方法中你可以 processing_date = context['execution_date'] - timedelta(days=3) .

    强制执行日期有一个延迟,感觉很不对。因为,根据气流的逻辑,当前运行的DAG的执行日期通常只有在追赶(bakcfilling)时才会有延迟。

        2
  •  1
  •   Daniel Huang    7 年前

    你可以用 TimeSensor 延迟DAG中任务的执行。我认为你不能改变实际情况 execution_date 除非你能把这种行为描述成一个老太婆。

    如果您希望它仅对计划的DAG运行的子集应用此延迟,可以使用 BranchPythonOperator 首先检查一下 执行日期 这是你想要滞后的一天。如果是,则使用传感器获取分支。否则,请不要使用它。

    或者,特别是如果您计划在多个DAG中使用此行为,则可以编写传感器的修改版本。它可能看起来像这样:

    def poke(self, context):
        if should_delay(context['execution_date']):
            self.log.info('Checking if the time (%s) has come', self.target_time)
            return timezone.utcnow().time() > self.target_time
        else:
            self.log.info('Not one of those days, just run')
            return True
    

    您可以在中引用现有时间传感器的代码 https://github.com/apache/incubator-airflow/blob/1.10.1/airflow/sensors/time_sensor.py#L38-L40 .

    推荐文章