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

