Airflow 2.8中DAG内运行FileSensor及相关问题排查
Airflow 2.8 FileSensor 常见问题解决
问题背景
在Airflow 2.8中使用自定义FileSensor生成器时遇到两个核心问题:
- 无法通过代码判断传感器是否找到文件,即便日志显示文件已被找到,预设的日志输出也未触发
- 无法正确提取
AIRFLOW_CTX_DAG_RUN_ID环境变量,导致日志中run_id显示为None
现有代码示例
通用FileSensor生成函数
def generic_human_feedback_sensor(step_name: str) -> FileSensor: prefix = "" run_id = str(os.getenv("AIRFLOW_CTX_DAG_RUN_ID")) if os.getenv("FILE_STORAGE_LAYER_ENV") == "local": pythonpath = os.getenv("PYTHONPATH") prefix = f"{pythonpath}/data/s3/" fs_connection_id = "fs_default" else: fs_connection_id = "aws_s3_log" definitions_validated_by_humans = prefix + f"input/{step_name}/{run_id}/**.json" return FileSensor( task_id=f"{step_name}_human_feedback", filepath=definitions_validated_by_humans, mode="poke", timeout=300, poke_interval=60, fs_conn_id=fs_connection_id, recursive=True, )
DAG中调用代码
@dag( dag_id="main", schedule=None, start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, tags=["main"], params={"pdf_path": Param("path to a pdf to run the DAG on.")}, ) def main(): ... definitions_sensor = generic_human_feedback_task(step_name="text_to_definitions") definitions_sensor.set_upstream(previous_task) if definitions_sensor is True: logging.error("File found")
日志信息
[2024-02-02, 18:35:16 CET] {base.py:83} INFO - Using connection ID 'fs_default' for task execution. [2024-02-02, 18:35:16 CET] {filesystem.py:66} INFO - Poking for file data/s3/input/<step name>/None/**.json [2024-02-02, 18:35:16 CET] {filesystem.py:71} INFO - Found File data/s3/input/<step name>/None/some_file.json last modified: 20240201184050
问题分析与解决方案
1. 传感器返回值判断无效的问题
- 原因:
definitions_sensor是FileSensor对象,而非布尔值。DAG定义阶段为静态解析,任务对象永远不会等于True。Airflow中传感器任务成功执行即代表找到目标文件,失败则代表超时未找到。 - 解决方法:通过任务依赖关系,在传感器任务成功后执行后续逻辑(比如打印日志),示例如下:
from airflow.operators.python import PythonOperator import logging def log_file_found(**context): logging.error("File found") @dag(...) def main(): ... definitions_sensor = generic_human_feedback_sensor(step_name="text_to_definitions") definitions_sensor.set_upstream(previous_task) # 传感器成功后执行日志打印任务 log_task = PythonOperator( task_id="log_file_found", python_callable=log_file_found, provide_context=True ) definitions_sensor >> log_task
2. 无法获取AIRFLOW_CTX_DAG_RUN_ID的问题
- 原因:
os.getenv("AIRFLOW_CTX_DAG_RUN_ID")在DAG解析阶段执行,此时Airflow尚未注入运行时环境变量,因此返回None。运行时环境变量仅在任务执行阶段可用。 - 解决方法:使用Airflow的Jinja模板语法,让
filepath在任务运行时动态渲染run_id:
def generic_human_feedback_sensor(step_name: str) -> FileSensor: prefix = "" if os.getenv("FILE_STORAGE_LAYER_ENV") == "local": pythonpath = os.getenv("PYTHONPATH") prefix = f"{pythonpath}/data/s3/" fs_connection_id = "fs_default" else: fs_connection_id = "aws_s3_log" # 使用Jinja模板{{ run_id }},运行时自动替换为当前DAG运行ID definitions_validated_by_humans = prefix + f"input/{step_name}/{{{{ run_id }}}}/**.json" return FileSensor( task_id=f"{step_name}_human_feedback", filepath=definitions_validated_by_humans, mode="poke", timeout=300, poke_interval=60, fs_conn_id=fs_connection_id, recursive=True, )
注意:Jinja模板需要用双大括号转义(
{{{ run_id }}}),避免DAG解析阶段被提前渲染。
内容的提问来源于stack exchange,提问作者Fares
相关产品推荐
相关产品推荐

