自定义Apache Airflow FileSensor:支持超时重试并创建缺失文件
自定义Airflow传感器实现文件检查与自动创建
问题概述
使用Airflow 2.8.1时,原方案通过FileSensor的on_retry_callback创建文件,但Airflow 2.2+版本中传感器超时后不会触发重试,导致回调无法执行,任务持续失败。需要实现自定义传感器,达成“检查文件存在→不存在则创建→直到文件存在后进入下一任务”的逻辑。
解决方案:扩展FileSensor实现自定义逻辑
直接继承FileSensor并重写poke方法,在每次检查环节主动创建缺失文件,无需依赖重试机制,确保每次探测都会尝试修复文件缺失问题。
自定义传感器代码
from airflow.sensors.filesystem import FileSensor from airflow.utils.decorators import apply_defaults import os class CreateFileOnMissingSensor(FileSensor): @apply_defaults def __init__(self, file_content: str = "Some random data\n", *args, **kwargs): super().__init__(*args, **kwargs) self.file_content = file_content def poke(self, context): # 调用父类方法检查文件是否存在 file_exists = super().poke(context) if not file_exists: # 确保文件目录存在,避免创建文件时报错 os.makedirs(os.path.dirname(self.filepath), exist_ok=True) # 创建并写入文件内容 with open(self.filepath, 'w') as f: f.write(self.file_content) # 再次检查文件是否创建成功 return super().poke(context) return file_exists
使用示例
from airflow import DAG from airflow.utils.dates import days_ago # 目标文件路径,注意需指定具体文件名,而非目录 FILE_PATH = "/mnt/c/path/to/file/data.txt" with DAG( dag_id="orchestrator-dag", schedule_interval='@daily', start_date=days_ago(1), catchup=False ) as testDag: check_create_file_task = CreateFileOnMissingSensor( task_id="check-and-create-file", filepath=FILE_PATH, poke_interval=5, timeout=30, # 可选:自定义文件内容,默认是"Some random data\n" file_content="Custom content for auto-created file\n" ) # 后续任务示例(可根据需求替换) # from airflow.operators.bash import BashOperator # next_task = BashOperator( # task_id="next-task", # bash_command="echo 'File is ready, proceed to next step'" # ) # check_create_file_task >> next_task
关键说明
- 重写
poke方法:每次传感器探测时主动检查文件,缺失则立即创建,绕过了Airflow 2.2+的传感器重试限制 - 目录预创建:通过
os.makedirs(..., exist_ok=True)确保文件所在目录存在,避免因目录缺失导致创建失败 - 灵活性:支持传入自定义文件内容,适配不同场景需求
- 兼容性:完美支持Airflow 2.8.1及所有2.2+版本
内容的提问来源于stack exchange,提问作者Lihka_nonem
相关产品推荐
相关产品推荐

