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

自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 21:26:26