如何获取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
相关产品推荐
相关产品推荐

