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

如何使用Airflow FileSensor检索最近修改的文件及场景实现

问题解答

需求理解判断

你的理解完全正确。原生FileSensor仅支持检测指定路径下是否存在匹配规则的文件,无法基于上次成功检测的时间戳筛选增量修改的文件,也无法批量传递符合要求的文件列表给下游任务,确实需要自定义Sensor组件实现需求。

上次成功poke时间戳的存储方案

不要将时间戳存储在自定义Sensor的实例变量中:Airflow的Sensor每次poke都是独立运行的实例,实例变量不会在两次poke之间持久化,worker重启、调度重触发都会导致数据丢失。
你可以通过XCom实现该数据的持久化存储,相关注意事项如下:

  • Airflow 1.10.15的XCom默认和任务实例、DAG运行的execution_date绑定,拉取历史时间戳时需要添加include_prior_dates=True参数,才能拉取到之前DAG运行实例中该任务推送的XCom数据
  • 对应场景的task_id就是你给自定义Sensor任务设置的id,比如命名为dir_sensor_check即可

筛选文件列表的传递方式

筛选得到的文件列表可以直接通过XCom传递给下游任务:你可以在自定义Sensor的poke方法返回True之前,分别把当前时间戳、筛选得到的文件列表push到XCom,下游任务开启provide_context=True后,通过对应task_id和key就能拉取到文件列表。
以下是自定义Sensor的参考实现:

from airflow.operators.sensors import BaseSensorOperator
from airflow.utils.decorators import apply_defaults
import os
from datetime import datetime

class MyDirSensor(BaseSensorOperator):
    @apply_defaults
    def __init__(self, dir_path, file_suffix='.csv', *args, **kwargs):
        super(MyDirSensor, self).__init__(*args, **kwargs)
        self.dir_path = dir_path
        self.file_suffix = file_suffix

    def poke(self, context):
        ti = context['ti']
        # 拉取上次成功检测的时间戳,首次运行默认取很早的时间兜底
        last_poke_ts = ti.xcom_pull(
            task_ids=self.task_id,
            key='last_success_poke_ts',
            include_prior_dates=True
        ) or datetime(2000, 1, 1).timestamp()
        
        current_ts = datetime.now().timestamp()
        matched_files = []
        # 遍历目录筛选增量修改的csv文件
        for file_name in os.listdir(self.dir_path):
            if not file_name.endswith(self.file_suffix):
                continue
            full_path = os.path.join(self.dir_path, file_name)
            file_mtime = os.path.getmtime(full_path)
            if file_mtime > last_poke_ts:
                matched_files.append(full_path)
        
        if matched_files:
            # 推送文件列表和新的时间戳到XCom
            ti.xcom_push(key='matched_csv_files', value=matched_files)
            ti.xcom_push(key='last_success_poke_ts', value=current_ts)
            return True
        return False

下游任务拉取文件列表的参考代码:

from airflow.operators.python_operator import PythonOperator

def run_prediction_service(**context):
    matched_files = context['ti'].xcom_pull(
        task_ids='dir_sensor_check',
        key='matched_csv_files'
    )
    # 此处添加读取文件、调用预测服务的逻辑

prediction_task = PythonOperator(
    task_id='run_prediction',
    python_callable=run_prediction_service,
    provide_context=True,
    dag=dag
)

更适配模拟场景的替代方案

你当前是单机测试场景,也可以不用自定义Sensor,直接用普通PythonOperator实现相同逻辑:将该任务的调度间隔设为3分钟,每次执行直接做时间戳比对、文件筛选的逻辑,无需处理Sensor的poke机制,实现更简单。
注意如果你的文件列表长度很大,超过Airflow默认48KB的XCom大小限制,可以修改airflow.cfg中的max_xcom_size参数调整上限,或者将文件列表写入本地临时文件,仅把临时文件路径通过XCom传递即可。

内容的提问来源于stack exchange,提问作者soumeng78

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 01:06:01