【发布时间】:2021-03-12 14:41:11
【问题描述】:
TL;DR 如何创建一个动态 EmailOperator,它从 XCom 属性的文件路径发送文件
大家好,
我正在使用 Apache Airflow 2.0.0.b2。我的问题是我的 DAG 创建了一个名称在运行时更改的文件。我想通过电子邮件将此文件发送出去,但在将动态文件名输入到我的 EmailOperator 时遇到问题。
我尝试过但失败的事情!:
-
为
files属性使用模板。files=["{{ ti.xcom_pull(key='OUTPUT_CSV') }}"],不幸的是,模板仅在运算符中的字段被标记为此类时才有效。
files不是 EmailOperator 上的模板字段 -
使用函数动态创建我的任务
def get_email_operator(?...): export_file_path = ti.xcom_pull(key='OUTPUT_CSV') email_subject = 'Some Subject' return EmailOperator( task_id="get_email_operator", to=['someemail@somedomain.net'], subject=email_subject, files=[export_file_path,], html_content='<br>', dag=current_dag) ..task3 >> get_email_operator() >> task4不幸的是,我似乎无法弄清楚如何将当前的
**kwargs或ti信息传递到我的函数调用中以获取当前文件路径。
编辑: Elad 在下面的回答让我朝着正确的方向前进。我唯一要做的就是在调用 op.execute() 时添加 kwargs
解决方案:
def get_email_operator(**kwargs):
export_file_path = kwargs['ti'].xcom_pull(key='OUTPUT_CSV')
email_subject = 'Termed Drivers - ' + date_string
op = EmailOperator(
task_id="get_email_operator",
to=['someemail@somedomain.net'],
subject=email_subject,
files=[export_file_path,],
html_content='<br>')
op.execute(kwargs)
【问题讨论】:
标签: directed-acyclic-graphs airflow