使用GCSToGCSOperator在同GCS存储桶内移动文件失败问题
GCS文件移动任务失败的解决方案
原代码核心问题
- 未定义变量
data:列表推导式中使用了未定义的data,实际应该用从XCom拉取的source_files - 函数内实例化Operator无效:Airflow的Operator必须在DAG构建阶段定义,运行时(PythonOperator函数内)动态创建的Operator不会被调度执行
修正后的完整代码
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.google.cloud.operators.gcs import GCSListObjectsOperator from google.cloud import storage from datetime import datetime def segregate_files(ti): source_bucket = "my-bucket" PROJECT_ID = "project" destination_bucket = "my-bucket" source_files = ti.xcom_pull(task_ids='list_file') print(f"拉取到的文件列表: {source_files}") # 修正变量未定义问题:用source_files替代data new_files = [file for file in source_files if 'abc' in file] print(f"筛选后的目标文件: {new_files}") if new_files: client = storage.Client(project=PROJECT_ID) source_bucket_obj = client.get_bucket(source_bucket) dest_bucket_obj = client.get_bucket(destination_bucket) for file_path in new_files: # 提取文件名,拼接目标路径 file_name = file_path.split('/')[-1] destination_path = f"test/{file_name}" # 执行移动操作:复制到目标路径后删除原文件 blob = source_bucket_obj.blob(file_path) dest_bucket_obj.copy_blob(blob, dest_bucket_obj, destination_path) blob.delete() print(f"已完成移动: {file_path} -> {destination_path}") with DAG( dag_id='gcs_move_abc_files', start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: # 修正delimiter参数:原参数错误,移除后正确获取new/下的所有文件 list_new_files = GCSListObjectsOperator( task_id='list_file', bucket='my-bucket', prefix='new/', do_xcom_push=True ) process_move_task = PythonOperator( task_id='process_and_move_files', python_callable=segregate_files, provide_context=True ) list_new_files >> process_move_task
关键修正说明
- 变量修正:将列表推导式中的
data替换为source_files,解决NameError - Operator使用方式修正:放弃在函数内动态创建
GCSToGCSOperator,改用Google Cloud Storage SDK直接在Python函数中执行文件移动(复制+删除),确保运行时可以执行操作 - GCSListObjectsOperator参数修正:移除错误的
delimiter='.txt',该参数用于分隔文件夹前缀,而非过滤文件后缀,若需筛选.txt文件可在Python函数中额外添加过滤逻辑
内容的提问来源于stack exchange,提问作者avinash reddy
相关产品推荐
相关产品推荐

