如何修改Airflow GCS传感器以检测指定前缀的新Avro文件更新?
解决GCS前缀路径下Avro文件更新的Airflow传感器检测问题
你当前使用的GCSObjectsWithPrefixExistenceSensor仅检测前缀路径下是否存在对象,无法识别文件的更新(比如旧文件被移除后上传的新文件),因此会出现传感器持续判定成功的问题。以下是针对性的解决方案:
自定义传感器实现文件更新检测
基于GCSObjectsWithPrefixExistenceSensor扩展,添加文件最后修改时间的对比逻辑,仅当检测到新上传的Avro文件时返回成功:
from airflow.providers.google.cloud.sensors.gcs import GCSObjectsWithPrefixExistenceSensor from airflow.utils.decorators import apply_defaults from google.cloud.storage import Client class GCSAvroFileUpdateSensor(GCSObjectsWithPrefixExistenceSensor): @apply_defaults def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.last_recorded_mod_time = None def poke(self, context): # 初始化GCS客户端 client = Client(credentials=self._get_credentials()) bucket = client.get_bucket(self.bucket) # 筛选前缀下的Avro文件,排除目录占位对象 avro_blobs = [ blob for blob in bucket.list_blobs(prefix=self.prefix) if blob.name.endswith('.avro') ] if not avro_blobs: return False # 获取当前唯一Avro文件的修改时间 current_mod_time = avro_blobs[0].updated # 首次检测:记录时间并返回成功 if not self.last_recorded_mod_time: self.last_recorded_mod_time = current_mod_time return True # 后续检测:仅当文件修改时间晚于上次记录时返回成功 elif current_mod_time > self.last_recorded_mod_time: self.last_recorded_mod_time = current_mod_time return True return False
使用自定义传感器替换原有代码
将原有的GCSObjectsWithPrefixExistenceSensor替换为自定义传感器即可:
gcs_avro_update_sensor = GCSAvroFileUpdateSensor( task_id='gcs_avro_file_update_sensor', bucket='my_bucket', prefix='mysql/abc/', google_cloud_conn_id='google_cloud_default', timeout=600, poke_interval=60, dag=dag )
逻辑说明
- 每次
poke时,传感器会扫描指定前缀下的所有Avro文件(确保仅处理目标文件类型) - 记录首次检测到的文件修改时间,后续每次检测都会对比当前文件的修改时间
- 只有当新文件上传(修改时间更新)时,传感器才会返回成功,触发后续任务处理新文件
内容的提问来源于stack exchange,提问作者Lukasz
相关产品推荐
相关产品推荐

