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

如何从Airflow的FileSensor中提取已检测文件的绝对路径?

问题描述

我正在构建Airflow DAG,使用FileSensor监听文件系统变化,对新增/修改的文件执行分析。监控路径包含Jinja模板和通配符/glob,希望检测到文件后,将其绝对路径传递给后续回调和任务,用于元数据对比判断是否需要处理。

核心问题:**如何从FileSensor中提取已找到文件的路径?**查看FileSensor源码发现它仅记录文件路径但未存储,想确认能否通过路径模板和上下文重构路径,无需额外文件系统查询。

目前想到两个临时方案,但希望找到更优方式:

  • 将数据路径模板直接传递给后续任务,期望在PythonOperator中自动生效;
  • 通过上下文/环境或BashOperator+echo强制渲染Jinja模板,再重新查询文件系统。

使用Airflow版本:2.8.0

简化配置代码

DAG初始化部分

# <DAG初始化代码>
    ...
    path_template: str = os.path.join(
        "/basepath/",
        "{{ data_interval_start.format('YYYYMMDD') }}",
        f"{source}.{{{{ data_interval_start.format('YYYYMMDD') }}}}*.csv.gz"
    )
    fs: FileSensor = FileSensor(
        task_id=f"{source}_data_sensor",
        filepath=path_template,
        poke_interval=int(duration(minutes=5).total_seconds()),
        timeout=int(duration(hours=1, minutes=30).total_seconds()),
        mode="reschedule",
        pre_execute=log_execution,
        on_success_callback=partial(log_found_file, path_template),
    )
    fs >> convert(source) >> analyze_data(source)

回调函数log_found_file

def log_found_file(data_path_template: str, ctx: Context) -> None:
    """Logs the discovery of a data file."""
    data_path = f(data_path_template, ctx)  # <<<<<<<<<<<<<<<<<< 需要此处的实现帮助
    stats: os.stat_result = os.stat(data_path)
    logger.success(
        f"Detected data file {data_path} "
        f"of size {stats.st_size}; "
        f"created on {stats.st_ctime}; "
        f"and last modified on {stats.st_mtime}."
    )

解决方案

1. 利用Airflow原生模板渲染+轻量glob匹配(最简方案)

Airflow提供render_template工具函数,可直接用任务上下文渲染模板字符串,结合glob匹配实际存在的文件(FileSensor已确认文件存在,查询开销极小)。修改log_found_file函数如下:

from airflow.utils.template import render_template
import glob
import os

def log_found_file(data_path_template: str, ctx: Context) -> None:
    """Logs the discovery of a data file."""
    # 渲染Jinja模板得到带通配符的路径
    rendered_path = render_template(data_path_template, ctx)
    # 匹配实际存在的文件
    matched_files = glob.glob(rendered_path)
    
    # 处理多文件场景(单文件场景取第一个即可)
    for data_path in matched_files:
        stats = os.stat(data_path)
        logger.success(
            f"Detected data file {data_path} "
            f"of size {stats.st_size}; "
            f"created on {stats.st_ctime}; "
            f"and last modified on {stats.st_mtime}."
        )

2. 自定义FileSensor子类,将路径存入XCom(跨任务传递最优方案)

如果需要后续任务直接复用路径,无需重复查询,可以自定义FileSensor子类,在检测到文件时将路径存入XCom:

from airflow.sensors.filesystem import FileSensor
from airflow.utils.template import render_template
from airflow.utils.xcom import XCom
import glob

class FileSensorWithPath(FileSensor):
    def poke(self, context):
        # 调用父类方法判断文件是否存在
        result = super().poke(context)
        if result:
            # 渲染模板并匹配文件
            rendered_path = render_template(self.filepath, context)
            matched_files = glob.glob(rendered_path)
            if matched_files:
                # 将第一个匹配的文件路径存入XCom(多文件可存列表)
                XCom.set(
                    key="detected_file_path",
                    value=matched_files[0],
                    task_id=self.task_id,
                    dag_id=self.dag_id,
                    execution_date=context["execution_date"]
                )
        return result

在DAG中替换为自定义传感器:

fs: FileSensorWithPath = FileSensorWithPath(
    task_id=f"{source}_data_sensor",
    filepath=path_template,
    # 其他参数保持不变
)

后续任务/回调直接从XCom读取路径:

def log_found_file(ctx: Context) -> None:
    data_path = ctx["ti"].xcom_pull(
        task_ids=f"{source}_data_sensor", 
        key="detected_file_path"
    )
    stats = os.stat(data_path)
    logger.success(
        f"Detected data file {data_path} "
        f"of size {stats.st_size}; "
        f"created on {stats.st_ctime}; "
        f"and last modified on {stats.st_mtime}."
    )

3. 优化临时方案1:让PythonOperator自动渲染模板

如果要直接传递模板给后续PythonOperator,可利用Airflow的templates_dict参数实现自动渲染,无需手动处理:

from airflow.operators.python import PythonOperator

def convert_task(source, detected_path):
    # 直接使用已渲染的文件路径
    print(f"Processing file: {detected_path}")
    # 你的转换逻辑...

convert_task_op = PythonOperator(
    task_id=f"convert_{source}",
    python_callable=convert_task,
    op_kwargs={
        "source": source,
        "detected_path": path_template
    },
    templates_dict={"detected_path": path_template},  # 标记需要渲染的参数
    provide_context=True,
)

总结
  • 仅回调使用路径:方案1最简,复用原生渲染+轻量查询;
  • 跨任务传递路径:方案2最优,一次查询多任务复用;
  • 延续临时方案思路:方案3利用Airflow原生能力减少代码量。

内容的提问来源于stack exchange,提问作者Dev-iL

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 07:16:13