You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 02:56:26