如何使用Airflow FileSensor检索最近修改的文件及场景实现
问题解答
需求理解判断
你的理解完全正确。原生FileSensor仅支持检测指定路径下是否存在匹配规则的文件,无法基于上次成功检测的时间戳筛选增量修改的文件,也无法批量传递符合要求的文件列表给下游任务,确实需要自定义Sensor组件实现需求。
上次成功poke时间戳的存储方案
不要将时间戳存储在自定义Sensor的实例变量中:Airflow的Sensor每次poke都是独立运行的实例,实例变量不会在两次poke之间持久化,worker重启、调度重触发都会导致数据丢失。
你可以通过XCom实现该数据的持久化存储,相关注意事项如下:
- Airflow 1.10.15的XCom默认和任务实例、DAG运行的execution_date绑定,拉取历史时间戳时需要添加
include_prior_dates=True参数,才能拉取到之前DAG运行实例中该任务推送的XCom数据 - 对应场景的task_id就是你给自定义Sensor任务设置的id,比如命名为
dir_sensor_check即可
筛选文件列表的传递方式
筛选得到的文件列表可以直接通过XCom传递给下游任务:你可以在自定义Sensor的poke方法返回True之前,分别把当前时间戳、筛选得到的文件列表push到XCom,下游任务开启provide_context=True后,通过对应task_id和key就能拉取到文件列表。
以下是自定义Sensor的参考实现:
from airflow.operators.sensors import BaseSensorOperator from airflow.utils.decorators import apply_defaults import os from datetime import datetime class MyDirSensor(BaseSensorOperator): @apply_defaults def __init__(self, dir_path, file_suffix='.csv', *args, **kwargs): super(MyDirSensor, self).__init__(*args, **kwargs) self.dir_path = dir_path self.file_suffix = file_suffix def poke(self, context): ti = context['ti'] # 拉取上次成功检测的时间戳,首次运行默认取很早的时间兜底 last_poke_ts = ti.xcom_pull( task_ids=self.task_id, key='last_success_poke_ts', include_prior_dates=True ) or datetime(2000, 1, 1).timestamp() current_ts = datetime.now().timestamp() matched_files = [] # 遍历目录筛选增量修改的csv文件 for file_name in os.listdir(self.dir_path): if not file_name.endswith(self.file_suffix): continue full_path = os.path.join(self.dir_path, file_name) file_mtime = os.path.getmtime(full_path) if file_mtime > last_poke_ts: matched_files.append(full_path) if matched_files: # 推送文件列表和新的时间戳到XCom ti.xcom_push(key='matched_csv_files', value=matched_files) ti.xcom_push(key='last_success_poke_ts', value=current_ts) return True return False
下游任务拉取文件列表的参考代码:
from airflow.operators.python_operator import PythonOperator def run_prediction_service(**context): matched_files = context['ti'].xcom_pull( task_ids='dir_sensor_check', key='matched_csv_files' ) # 此处添加读取文件、调用预测服务的逻辑 prediction_task = PythonOperator( task_id='run_prediction', python_callable=run_prediction_service, provide_context=True, dag=dag )
更适配模拟场景的替代方案
你当前是单机测试场景,也可以不用自定义Sensor,直接用普通PythonOperator实现相同逻辑:将该任务的调度间隔设为3分钟,每次执行直接做时间戳比对、文件筛选的逻辑,无需处理Sensor的poke机制,实现更简单。
注意如果你的文件列表长度很大,超过Airflow默认48KB的XCom大小限制,可以修改airflow.cfg中的max_xcom_size参数调整上限,或者将文件列表写入本地临时文件,仅把临时文件路径通过XCom传递即可。
内容的提问来源于stack exchange,提问作者soumeng78
相关产品推荐
相关产品推荐

