【问题标题】:In Airflow how can I pass parameters using context to on_success_callback function handler?在 Airflow 中,如何使用上下文将参数传递给 on_success_callback 函数处理程序?
【发布时间】:2021-02-09 15:56:37
【问题描述】:

在 Airflow 中,如何使用上下文将参数传递给 on_success_callback 函数处理程序?

这是我的测试代码:

import airflow
from airflow import DAG
from airflow.operators import MSTeamsWebhookOperator
from airflow.operators.bash_operator import BashOperator
from airflow.operators.dummy_operator import DummyOperator
from datetime import datetime
from transaction_analytics import helpers
from airflow.utils.helpers import chain

# Parameters & variables
schedule_interval        = "0 20 * * *"

def _task_success_callback(context):
    dagid         = context["task_instance"].dag_id
    duration      = context["task_instance"].duration
    executiondate = context["execution_date"]
    logurl        = context["task_instance"].log_url.replace("localhost", "agbqhsbldd017v.agb.rbxd.ds")# workaround until we config airflow
    pp1           = context["params"].param1 
    #pp1          = "{{ params.param1 }}"
    ms_teams_op = MSTeamsWebhookOperator(
        task_id="success_notification",
        http_conn_id="msteams_airflow",
        message="DAG {ppram1} `{dag}`  finished successfully!".format(dag=context["task_instance"].dag_id, ppram1=pp1),
        subtitle="Execution Date = {p1}, Duration = {p2}".format(p1=executiondate,p2=duration),
        button_text = "View log",
        button_url = "{log}".format(log=logurl),
        theme_color="00FF00"#,
        #proxy= "http://10.72.128.202:3128"
    )
    ms_teams_op.execute(context)

    main_dag = DAG('test_foley',
            schedule_interval=schedule_interval,
            description='Test foley',
            start_date=datetime(2020, 4, 19),
            default_args=None,
            max_active_runs=2,
            default_view='graph',     # Default view graph
            #orientation='TB', # Top-Bottom graph
            on_success_callback=_task_success_callback,
            #on_failure_callback=outer_task_failure_callback,
            catchup=False, # Do not catchup, run only latest
            params={
              "param1": "value1",
              "param2": "value2"
            }
           )

################################### START ######################################
dag_chain = []

start = DummyOperator(task_id='start', retries = 3, dag=main_dag)
dag_chain.append(start)

step1 = BashOperator(
    task_id='step1',
    bash_command='pwd',
    dag=main_dag,
)
dag_chain.append(step1)

step2 = BashOperator(
    task_id='step2',
    bash_command='exit 0',
    dag=main_dag,
)
dag_chain.append(step2)


end = DummyOperator(task_id='end', dag=main_dag)
dag_chain.append(end)

chain(*dag_chain)

我有一个处理成功的事件处理函数_task_success_callback。 在 DAG 中,我有 on_success_callback=_task_success_callback 来捕获该事件。

它可以工作......但现在我需要将一些参数传递给 _task_success_callback。 最好的方法是什么?

当该函数接收上下文时,我尝试在 DAG 中创建参数,如您所见:

        params={
          "param1": "value1",
          "param2": "value2"
        }

但似乎我无法访问它们?

我的问题是:

  1. 访问参数我做错了什么?
  2. 有没有更好的参数传递方式?

注意:我看到这个类似的问题How to pass parameters to Airflow on_success_callback and on_failure_callback 有一个答案......并且有效。但我正在寻找的是使用上下文来传递参数....

【问题讨论】:

  • 是否有特定原因导致您希望使用上下文而您的链接答案不起作用?
  • 只是为了有更简洁的代码。我有许多 DAG,每个 DAG 都通知团队在 MsTeamsWebHook 运算符中具有不同的值。使用当前的解决方案,我必须将 DAG 连接到 2 个函数(成功和失败),并将这些函数连接到库中的公共函数。我有 32 个 DAG,这意味着还要创建 64 个函数(每个 DAG 2 个),唯一的区别是参数(http_conn_id、消息、标题)。如果我能够通过上下文传递参数,我将有 32 个 DAG 直接调用带参数的通用函数,让我免于使用那些丑陋的 64 个函数......
  • 必须是context["params"]['param1'] 而不是context["params"].param1

标签: python airflow


【解决方案1】:

回想一下,Airflow 流程文件只是 Python,如果您在解析过程中没有引入太多开销(由于 Airflow 经常解析文件,而且开销会增加),您可以使用 Python 能做的一切。特别是对于您的情况,我建议为您的回调返回一个嵌套函数 (closure):

将其放在与 Airflow 进程相邻的文件中,例如 on_callbacks.py

def success_ms_teams(param_1, param_2):

    def callback_func(context):
        print(f"param_1: {param_1}")
        print(f"param_2: {param_2}")
        # ... trimmed for brevity ...#
        ms_teams_op.execute(context)

    return callback_func

然后在您的流程中,您可以这样做:

from airflow import models

from on_callbacks import success_ms_teams

with models.DAG(
    ...
    on_success_callback=success_ms_teams(
        "value1", # These values become the
        "value2", # `param_1` and `param_2` 
    )
) as dag:
    ...

【讨论】:

  • 这会产生以下错误`python_callable param must be callable`
  • 您传递给python_callable 的值不可调用。尝试使用您的代码和日志提出问题。
【解决方案2】:

您可以创建一个任务,其唯一目的是通过xcoms 推送配置设置。您可以通过context 拉取配置,因为task_instance 对象包含在context 中。

def push_configuration(ti, params):
    ti.xcom_push(key='conn_id', value=params)

def _task_success_callback(context):
    ti = context.get('ti') 
    params = ti.xcom_pull(key='params', task_ids='Settings')
    ...


step0 = PythonOperator(
        task_id='Settings',
        python_callable=push_configuration,
        op_kwargs={'params': params})

step1 = BashOperator(
        task_id='step1',
        bash_command='pwd',
        on_success_callback=_task_success_callback)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-11-30
    • 2011-05-02
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多