如何使用Composer DAG递归读取GCS存储桶文件名?是否可行?
当然可行!用Composer DAG递归遍历GCS存储桶文件的实现方案
完全可以通过Airflow(Composer基于的核心框架)实现递归读取GCS存储桶下所有层级的文件,还能将存储桶名和每个文件的完整路径存入指定变量。下面给你具体的实现思路和代码示例:
核心思路
借助Airflow的GoogleCloudStorageHook与GCS交互,通过list方法的recursive=True参数,就能一键递归获取存储桶内所有文件的相对路径。之后可以将这些路径与存储桶名拼接,拆分出你需要的bucketname和filepath,并通过Airflow变量或XCom进行存储传递。
完整DAG代码示例
from airflow import DAG from airflow.providers.google.cloud.hooks.gcs import GoogleCloudStorageHook from airflow.operators.python import PythonOperator from airflow.utils.dates import days_ago from airflow.models import Variable default_args = { 'owner': 'airflow', 'start_date': days_ago(1), } def recursive_list_gcs_files(**context): # 定义目标存储桶名 target_bucket = "static" # 初始化GCS Hook(需确保Composer已配置好GCP连接) gcs_hook = GoogleCloudStorageHook(gcp_conn_id='google_cloud_default') # 递归获取存储桶内所有文件的相对路径 relative_file_paths = gcs_hook.list(bucket_name=target_bucket, recursive=True) # 遍历处理每个文件路径,存储所需变量 all_file_info = [] for rel_path in relative_file_paths: # 拼接完整文件路径 full_file_path = f"{target_bucket}/{rel_path}" # 组装需要的变量结构 file_info = { "bucketname": target_bucket, "filepath": full_file_path } all_file_info.append(file_info) # 可选:通过XCom传递给后续任务(适合任务间临时数据传递) context['ti'].xcom_push(key=f"file_{rel_path.replace('/', '_')}", value=file_info) # 将所有文件信息存入Airflow全局变量(适合长期保存) Variable.set("gcs_static_bucket_files", all_file_info, serialize_json=True) print(f"成功递归读取 {len(relative_file_paths)} 个文件,信息已存入变量") with DAG( 'gcs_recursive_file_scanner', default_args=default_args, description='递归读取GCS存储桶所有文件并存储路径信息', schedule_interval=None, catchup=False, ) as dag: scan_files_task = PythonOperator( task_id='recursive_scan_gcs_files', python_callable=recursive_list_gcs_files, provide_context=True, ) scan_files_task
关键细节说明
- 递归遍历的核心:
gcs_hook.list()方法中的recursive=True是实现子文件夹遍历的关键,缺省该参数只会读取存储桶根目录的文件。 - 变量存储选择:
- 如果需要长期保存文件路径信息,使用Airflow的
Variable更合适,注意用serialize_json=True来处理列表类型数据。 - 如果只是在DAG的任务间临时传递数据,用XCom更轻量,无需持久化到数据库。
- 如果需要长期保存文件路径信息,使用Airflow的
- 权限要求:确保Composer的服务账号拥有目标GCS存储桶的
storage.objects.list权限,否则会出现访问被拒绝的错误。
结果验证
执行这个DAG后,你可以通过Airflow UI的「Variables」页面查看gcs_static_bucket_files变量,里面会包含所有文件的信息,格式如下:
[ {"bucketname": "static", "filepath": "static/folder1/subfolder1/file1.json"}, {"bucketname": "static", "filepath": "static/folder1/subfolder2/file2.json"}, {"bucketname": "static", "filepath": "static/folder1/subfolder3/file3.json"}, {"bucketname": "static", "filepath": "static/folder1/subfolder3/file4.json"} ]
内容的提问来源于stack exchange,提问作者Shilaba Roul
相关产品推荐
相关产品推荐

