使用Airflow按匹配模式删除GCS对象时提示404不存在问题咨询
问题成因
GCSDeleteObjectsOperator的objects参数仅接收精确的GCS对象完整路径列表,原生不支持通配符(*)匹配逻辑。你传入带通配符的字符串后,Operator会直接把该字符串当作完整对象名调用GCS删除接口,GCS侧找不到名称完全一致的对象,就会返回404错误。
正确实现方案
GCS本身的列表API仅支持前缀匹配,不支持任意位置的通配符检索,要实现通配规则删除需要分两步走:
- 先调用GCS列表接口拉取固定前缀下的所有对象,再在本地通过通配符规则过滤出需要删除的目标对象完整路径
- 将过滤得到的精确路径列表传入
GCSDeleteObjectsOperator执行删除
可参考以下实现代码:
from airflow.providers.google.cloud.hooks.gcs import GCSHook from airflow.providers.google.cloud.operators.gcs import GCSDeleteObjectsOperator from airflow.operators.python import PythonOperator from fnmatch import fnmatch from airflow.models import Variable BUCKET_NAME = Variable.get("your_gcs_bucket_name") GCP_CONN_ID = "your_gcp_connection_id" DELETE_PATTERN = "test_delete/*/*/*/alpha/data-1-2123-*.json" def filter_target_objects(**context): gcs_hook = GCSHook(gcp_conn_id=GCP_CONN_ID) # 拉取固定前缀下的所有对象,减少本地过滤的数据量 all_objects = gcs_hook.list( bucket_name=BUCKET_NAME, prefix="test_delete/", delimiter="" ) # 按通配规则过滤目标对象 matched_objects = [obj for obj in all_objects if fnmatch(obj, DELETE_PATTERN)] if not matched_objects: raise RuntimeError("No objects matched the delete pattern, abort task") return matched_objects get_matched_objects = PythonOperator( task_id="get_matched_objects", python_callable=filter_target_objects, provide_context=True ) delete_data = GCSDeleteObjectsOperator( bucket_name=BUCKET_NAME, task_id="delete_data", objects="{{ ti.xcom_pull(task_ids='get_matched_objects') }}" ) get_matched_objects >> delete_data
注意事项
- 如果待匹配的对象数量较多,调用
gcs_hook.list时需要开启分页配置,避免漏取对象 - 不要直接给
GCSDeleteObjectsOperator传prefix参数删除整个前缀下的内容,会误删不符合路径规则的非目标对象
内容的提问来源于stack exchange,提问作者Daljeet Singh
相关产品推荐
相关产品推荐

