Airflow 2.2 Email Operator读取XCom文件列表失败求助
Airflow 2.2 EmailOperator动态发送多文件的问题解决
问题根源
你遇到的FileNotFoundError: [Errno 2] No such file or directory: '[''错误,本质是Jinja模板将XCom返回的列表渲染成了字符串形式,而EmailOperator的files参数期望接收Python列表对象,而非"['/path/file1', '/path/file2']"这种字符串格式的列表。EmailOperator会把整个字符串当作单个文件名去查找,自然会找不到文件。
解决方案
直接使用PythonOperator封装邮件发送逻辑,在Python代码中直接获取XCom返回的列表,调用Airflow内置的send_email函数完成发送,这样能直接处理列表对象,避免Jinja渲染带来的类型问题。同时优化路径拼接逻辑,避免手动拼接斜杠的潜在错误。
修改后的完整代码
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.email import send_email import os from datetime import datetime, timedelta default_args = { "owner": 'TEST', "depends_on_past": False, "email_on_failure": False, "email_on_retry": False, "retries": 0, } with DAG( dag_id="test9_email_operator_dag", default_args=default_args, start_date=datetime(2022, 9, 14), end_date=datetime(2022, 9, 15), catchup=True, max_active_runs=1, schedule_interval="0 12 * * *", # Runs every day @ 8AM EST ) as dag: def collect_files(local_temp_folder): print("local folder files => ", os.listdir(local_temp_folder)) files_list = [] for file in os.listdir(local_temp_folder): # 用os.path.join自动处理路径斜杠,避免手动拼接错误 file_path = os.path.join(local_temp_folder, file) files_list.append(file_path) print("files_list => ", files_list) return files_list collect_files_task = PythonOperator( task_id='collect_files', python_callable=collect_files, op_kwargs={'local_temp_folder': "/usr/local/airflow/dags/temp_dir/"}, do_xcom_push=True, dag=dag) def send_email_with_files(**context): # 从上下文直接获取XCom中的文件列表 files_list = context['task_instance'].xcom_pull(task_ids='collect_files') send_email( to='test@gmail.com', subject='Test Email op Notification', html_content='Test email op notification email. ', files=files_list ) send_email_task = PythonOperator( task_id='send_email', python_callable=send_email_with_files, provide_context=True, dag=dag) collect_files_task >> send_email_task
内容的提问来源于stack exchange,提问作者user13274071
相关产品推荐
相关产品推荐

