【问题标题】:Apache Airflow 2.0.0.b2 - Dynamic EmailOperator [files] propertyApache Airflow 2.0.0.b2 - 动态 EmailOperator [files] 属性
【发布时间】:2021-03-12 14:41:11
【问题描述】:

TL;DR 如何创建一个动态 EmailOperator,它从 XCom 属性的文件路径发送文件

大家好,

我正在使用 Apache Airflow 2.0.0.b2。我的问题是我的 DAG 创建了一个名称在运行时更改的文件。我想通过电子邮件将此文件发送出去,但在将动态文件名输入到我的 EmailOperator 时遇到问题。

我尝试过但失败的事情!:

  1. files 属性使用模板。

    files=["{{ ti.xcom_pull(key='OUTPUT_CSV') }}"],
    

    不幸的是,模板仅在运算符中的字段被标记为此类时才有效。 files 不是 EmailOperator 上的模板字段

  2. 使用函数动态创建我的任务

     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
    

    不幸的是,我似乎无法弄清楚如何将当前的**kwargsti 信息传递到我的函数调用中以获取当前文件路径。

编辑: 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


    【解决方案1】:

    由于上周合并了 PR,文件将在 Airflow 2 中进行模板化。

    但是您无需等待,您可以使用您自己的自定义操作符来包装当前操作符,指定模板化字段列表。

    喜欢:

    class MyEmailOperator(EmailOperator):
         template_fields = ('to', 'subject', 'html_content', 'files')
    

    然后你可以在你的代码中使用MyEmailOperatorfiles 将被模板化。

    你也可以通过使用包裹 EmailOperator 的 PythonOperator 来解决这个问题:

    def get_email_operator(**context):
        xcom = context['ti'].xcom_pull(task_ids='OUTPUT_CSV')
        email_subject = 'Some Subject'
        op = EmailOperator(
            task_id="get_email_operator",
            to=['someemail@somedomain.net'],
            subject=email_subject,
            files=[xcom,],
            html_content='<br>')
        op.execute(context)
    
    python = PythonOperator(
        task_id='archive_s3_file',
        dag=dag,
        python_callable=get_email_operator,
        provide_context=True
    )
    
    ..task3 >> python >> task4
    

    【讨论】:

    • 非常感谢。我选择了第二种解决方案,因为担心在经验不足的情况下错误地设计了我的管道。一旦我将上下文传递给 EmailOperator 的执行方法,它就像一个魅力。
    • @BlackDynamite NP。我建议您至少尝试第一个解决方案。它更简单,当修复程序将在下一版本中发布时,您可以轻松删除自定义操作符并将其替换为 Airflow 操作符。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-10-21
    • 1970-01-01
    • 2018-09-03
    • 2015-07-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多