通过Airflow从Azure ADLS向GCS复制100+文件遇问题求助
从Azure ADLS批量复制文件到GCS的替代方案
针对AzureBlobStorageToGCSOperator批量复制(100+文件)用通配符报“blob未找到”的问题,给你几个可行的替代方案:
方案一:分批次枚举+单个文件复制
先枚举所有目标文件,再逐个复制,避开通配符批量匹配的问题:
- 用
AzureBlobStorageListOperator列出ADLS中符合条件的所有blob路径,结果存入XCom - 借助
TaskGroup或Python循环,遍历每个路径调用AzureBlobStorageToGCSOperator单独复制 - 示例代码:
from airflow import DAG from airflow.providers.microsoft.azure.operators.wasb_list import AzureBlobStorageListOperator from airflow.providers.microsoft.azure.transfers.azure_blob_to_gcs import AzureBlobStorageToGCSOperator from airflow.operators.python import PythonOperator from airflow.utils.task_group import TaskGroup from datetime import datetime def copy_files(**context): blob_list = context['ti'].xcom_pull(task_ids='list_blobs') with TaskGroup("copy_batch_files") as batch_group: for blob in blob_list: # 替换路径中的特殊字符作为task_id safe_task_id = f"copy_{blob.replace('/', '_').replace('.', '_')}" AzureBlobStorageToGCSOperator( task_id=safe_task_id, wasb_conn_id="azure_adls_conn", container_name="source-container", blob_name=blob, gcp_conn_id="gcs_conn", dest_gcs="gs://target-bucket/", dest_path=f"your/dest/path/{blob}" ) return batch_group with DAG('adls_to_gcs_batch_copy', start_date=datetime(2024,1,1), schedule_interval=None) as dag: list_blobs = AzureBlobStorageListOperator( task_id='list_blobs', wasb_conn_id="azure_adls_conn", container_name="source-container", prefix="your/source/path/", delimiter="/" ) init_copy = PythonOperator( task_id='initiate_batch_copy', python_callable=copy_files, provide_context=True ) list_blobs >> init_copy
方案二:用Azure Data Factory(ADF)完成批量复制
ADF对批量文件复制的支持更稳定,自带重试、断点续传机制:
- 在ADF中创建复制活动,源选Azure Data Lake Storage Gen2,目标选Google Cloud Storage,配置好路径和批量规则
- 在Airflow里用
AzureDataFactoryRunPipelineOperator触发ADF的复制管道,完成后再执行后续业务逻辑
方案三:自定义Python脚本调用SDK实现复制
直接用云服务商的SDK写复制逻辑,完全控制流程:
- 用Azure Storage SDK枚举ADLS中的文件,用GCS SDK逐个上传,还能自定义重试、日志和异常处理
- 示例代码片段:
from azure.storage.blob import BlobServiceClient from google.cloud import storage import os from tenacity import retry, stop_after_attempt, wait_exponential @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) def upload_to_gcs(bucket, blob_name, data): gcs_blob = bucket.blob(blob_name) gcs_blob.upload_from_string(data) def batch_copy(): # 从环境变量取连接信息,避免硬编码 azure_conn_str = os.getenv("AZURE_STORAGE_CONN_STR") blob_service_client = BlobServiceClient.from_connection_string(azure_conn_str) container_client = blob_service_client.get_container_client("source-container") gcs_client = storage.Client() bucket = gcs_client.get_bucket("target-bucket") # 枚举所有要复制的文件 for blob in container_client.list_blobs(prefix="your/source/path/"): if not blob.name.endswith("/"): # 跳过文件夹 blob_data = container_client.download_blob(blob.name).readall() dest_path = f"your/dest/path/{blob.name}" upload_to_gcs(bucket, dest_path, blob_data)
- 把这段逻辑放进Airflow的
PythonOperator里执行即可
内容的提问来源于stack exchange,提问作者Khilesh Chauhan
相关产品推荐
相关产品推荐

