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

如何获取Apache Airflow FileSensor检测到的.xlsx文件名?

获取FileSensor检测到的文件名(通配符场景)

当你用*.xlsx作为FileSensor的filepath参数时,默认的FileSensor只会返回是否存在匹配文件,不会记录具体文件名。你可以通过以下两种方式获取检测到的文件名,再传入转换函数:

方法1:自定义FileSensor推送文件名到XCom

继承原生FileSensor,重写poke方法,找到匹配文件后将文件名存入XCom,供后续任务调用。

from airflow.sensors.filesystem import FileSensor
from airflow.hooks.filesystem import FileSystemHook
from airflow.utils.decorators import apply_defaults
import os

class CustomFileSensor(FileSensor):
    @apply_defaults
    def __init__(self, **kwargs):
        super().__init__(**kwargs)

    def poke(self, context):
        hook = FileSystemHook(self.fs_conn_id)
        # 获取所有匹配的xlsx文件
        matching_files = hook.get_files(self.filepath)
        if matching_files:
            # 取第一个匹配文件,可根据需求改为取最新修改的文件
            detected_file = os.path.basename(matching_files[0])
            # 将文件名推送到XCom
            context['ti'].xcom_push(key='detected_filename', value=detected_file)
            return True
        return False

# 使用自定义Sensor替代原生FileSensor
file_sensor_task = CustomFileSensor(
    task_id='waiting_for_file',
    filepath='*.xlsx',
    fs_conn_id="depot",
    poke_interval=30,
    retries=999,
    retry_delay=timedelta(seconds=10), 
    timeout=60 * 5,
    mode="reschedule",
    soft_fail=False,
)

后续转换任务通过XCom获取文件名:

from airflow.operators.python import PythonOperator

def xlsx_to_sqlite(filename):
    # 你的xlsx转SQLite逻辑
    print(f"正在处理文件: {filename}")

convert_task = PythonOperator(
    task_id='convert_xlsx_to_sqlite',
    python_callable=xlsx_to_sqlite,
    op_args=["{{ ti.xcom_pull(task_ids='waiting_for_file', key='detected_filename') }}"],
    dag=dag,
)

file_sensor_task >> convert_task

方法2:用PythonOperator实现检测+获取文件名

直接用PythonOperator完成文件轮询检测,同时返回文件名,省去自定义Sensor的步骤:

from airflow.operators.python import PythonOperator
from airflow.hooks.filesystem import FileSystemHook
import time
import os

def wait_for_xlsx_file(fs_conn_id, filepath, poke_interval, timeout):
    hook = FileSystemHook(fs_conn_id)
    start_time = time.time()
    while time.time() - start_time < timeout:
        matching_files = hook.get_files(filepath)
        if matching_files:
            # 按文件修改时间排序,取最新的文件
            matching_files.sort(key=lambda x: hook.get_mod_time(x), reverse=True)
            return os.path.basename(matching_files[0])
        time.sleep(poke_interval)
    raise TimeoutError(f"超时未找到匹配的xlsx文件: {filepath}")

# 文件检测任务
file_check_task = PythonOperator(
    task_id='waiting_for_file',
    python_callable=wait_for_xlsx_file,
    op_kwargs={
        'fs_conn_id': 'depot',
        'filepath': '*.xlsx',
        'poke_interval': 30,
        'timeout': 60*5
    },
    dag=dag,
)

# 转换任务
def xlsx_to_sqlite(filename):
    # 你的xlsx转SQLite逻辑
    print(f"正在处理文件: {filename}")

convert_task = PythonOperator(
    task_id='convert_xlsx_to_sqlite',
    python_callable=xlsx_to_sqlite,
    op_args=["{{ ti.xcom_pull(task_ids='waiting_for_file') }}"],
    dag=dag,
)

file_check_task >> convert_task

注意事项

  • 如果存在多个匹配的.xlsx文件,上述代码默认取第一个或最新修改的文件,你可以根据业务需求调整文件筛选逻辑。
  • 确保fs_conn_id="depot"对应的连接配置正确,能正常访问目标存储路径。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 09:35:18