You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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更轻量,无需持久化到数据库。
  • 权限要求:确保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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 08:41:36