【发布时间】: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"
}
但似乎我无法访问它们?
我的问题是:
- 访问参数我做错了什么?
- 有没有更好的参数传递方式?
注意:我看到这个类似的问题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?