【问题标题】:Can I programmatically determine if an Airflow DAG was scheduled or manually triggered?我能否以编程方式确定 Airflow DAG 是计划触发还是手动触发?
【发布时间】:2021-11-01 05:22:35
【问题描述】:

我想创建一个 sn-p,它根据 DAG 是否已安排或是否手动触发来传递正确的日期。 DAG 每月运行一次。 DAG 根据上个月的数据生成报告(A SQL 查询)。

如果我按计划运行 DAG,我可以使用以下 jinja sn-p 获取上个月:

execution_date.month

鉴于 DAG 被安排在上一个时期(上个月)结束时, execution_date 将正确返回上个月。但是在手动运行时,这将返回当前月份(执行日期将是手动触发的日期)。

我想编写一个简单的宏来处理这种情况。但是,我找不到以编程方式查询 DAG 是否以编程方式触发的好方法。我能想到的最好办法是从数据库中获取run_id(通过创建一个具有数据库会话的宏),检查run_id 是否包含单词manual。有没有更好的方法来解决这个问题?

【问题讨论】:

    标签: airflow


    【解决方案1】:

    tl;dr:您可以通过DagRun.external_trigger 确定这一点。


    我注意到在树视图中,有一个围绕已安排但非手动运行的大纲。那是因为后者在 CSS 中应用了stroke-opacity: 0;

    在 repo 中搜索这个,我发现了 Airflow devs detect manual runs(5 岁的行,所以也应该在旧版本中工作):

    .style("stroke-opacity", function(d) {return d.external_trigger ? "0": "1"})
    

    搜索external_trigger 会将我们带到DagRun definition

    因此,例如,如果您使用的是 Python 回调,则可以有这样的东西(可以在 DAG 或单独的文件中定义):

    def my_fun(context):
        if context.get('dag_run').external_trigger:
            print('manual run')
        else:
            print('scheduled run')
    

    并在您的Operator 中设置如下参数:

    t1 = BashOperator(
        task_id='print_date',
        bash_command='date',
        on_failure_callback=my_fun,
        dag=dag,
    )
    

    我已经测试过类似的东西并且它有效。

    我认为您也可以执行 if if {{ dag_run.external_trigger }}: 之类的操作 - 但我尚未对此进行测试,并且我相信它仅适用于该 DAG 的文件。

    【讨论】:

      【解决方案2】:

      目前没有直接的 DAG 属性来识别手动运行。 要获取此信息,您需要检查您提到的run_id

      但是,有一个专用的宏可以获取run_id。您不必自己从数据库中获取它。 这是一个如何使用它的示例:

          def some_task_py(**context):
              run_id = context['templates_dict']['run_id']
              is_manual = run_id.startswith('manual__')
              is_scheduled = run_id.startswith('scheduled__')
      
      
          some_task = PythonOperator(
                      task_id = 'some_task',
                      dag=dag,
                      templates_dict = {'run_id': '{{ run_id }}'},
                      python_callable = some_task_py,
                      provide_context = True)
      

      【讨论】:

      • 耻辱,确实是我所害怕的。 run_id 宏变量确实是一个很好的提示,我忘记了它的存在。
      • @Blokje5 出于好奇,您觉得这种方法的缺陷是什么?
      • @Donentolon 理想情况下,Airflow 会提供一些东西。如果下线有 10 个版本,他们将决定更改 run_id 生成机制,气流提供的机制仍然可以工作,而这会中断。我并不是说它会发生,但在这种情况下,我和 Airflow 之间没有合同。
      【解决方案3】:

      根据@Donentolon 的回答,我已经通过从kwargs 中的kwargs(我的PythonOperator)获取dag_run 来确定DAG 是手动触发还是计划触发:

      def my_python_callable(**kwargs):
          dag_run = kwargs.get("dag_run")
      
          if dag_run.external_trigger:
              logger.info("DAG triggered manually, skipping this operator")
              return True
      
          # my operator logic for scheduled run
      

      【讨论】:

        【解决方案4】:

        我需要能够检测是否计划或触发了某些事情,包括使用气流任务测试(或旧的气流测试)从命令行运行。

        我的任务参数是**kwargs,而不是像def some_task_py(**context) 这样的上下文,所以我的示例使用kwargs

        如果您从命令行运行,我相信kwargs['dag_run'] 将是None,而kwargs['templates_dict']['run_id'] 将不存在。

        我已经测试过,这应该可以从命令行运行或从计划或手动触发的 Web 服务器运行:

        if kwargs['dag_run'] == None or (kwargs['dag_run'] != None and kwargs['dag_run'].external_trigger):
            print("This is an external run.  Mark it as such")
        else:
            print("This is a scheduled run")
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2015-09-01
          • 2017-01-28
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2021-12-08
          相关资源
          最近更新 更多