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

通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 20:12:36