如何通过DAG代码在处理前检测文件大小,避免0字节文件处理
解决方案:在Airflow DAG中前置检测GCS文件大小
方法1:编写自定义文件大小检测任务
在DAG的起始阶段添加一个PythonOperator任务,通过GCP存储钩子获取文件元数据,判断大小是否为0,若为0则抛出异常终止流程。
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.google.cloud.hooks.gcs import GoogleCloudStorageHook from datetime import datetime def check_gcs_file_size(bucket_name, file_path): hook = GoogleCloudStorageHook(gcp_conn_id='your_gcp_connection_id') file_metadata = hook.get_metadata(bucket_name, file_path) file_size = file_metadata.get('size', 0) if file_size == 0: raise ValueError(f"文件 {file_path} 为0字节,终止处理") print(f"文件 {file_path} 大小 {file_size} 字节,继续执行后续任务") with DAG( dag_id='gcs_file_processing_dag', start_date=datetime(2024, 1, 1), schedule_interval='@daily', catchup=False ) as dag: check_file_task = PythonOperator( task_id='check_gcs_file_size', python_callable=check_gcs_file_size, op_kwargs={ 'bucket_name': 'your-target-bucket', 'file_path': 'data/input/file.csv' } ) # 后续处理任务示例 # process_task = ... check_file_task >> process_task
方法2:自定义Sensor过滤非0字节文件
如果是监听GCS桶中的新文件,可基于GoogleCloudStorageObjectSensor扩展,只在检测到非0字节文件时触发后续流程。
示例代码:
from airflow.providers.google.cloud.sensors.gcs import GoogleCloudStorageObjectSensor class GCSNonZeroFileSensor(GoogleCloudStorageObjectSensor): def poke(self, context): hook = GoogleCloudStorageHook(gcp_conn_id=self.gcp_conn_id) if hook.exists(self.bucket, self.object): file_metadata = hook.get_metadata(self.bucket, self.object) return file_metadata.get('size', 0) > 0 return False # 在DAG中使用自定义Sensor file_sensor = GCSNonZeroFileSensor( task_id='wait_for_valid_file', bucket_name='your-target-bucket', object='data/input/*.csv', gcp_conn_id='your_gcp_connection_id', poke_interval=60 )
关键注意事项
- 确保Airflow中已正确配置GCP连接(
gcp_conn_id对应后台配置的连接ID) - 处理批量文件时,可遍历桶内文件列表,批量检测大小后将有效文件路径传入后续任务
- 可配置
on_failure_callback在检测到0字节文件时发送告警通知
内容的提问来源于stack exchange,提问作者Ayush331
相关产品推荐
相关产品推荐

