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

请求协助:在Airflow中创建每日定时DAG实现S3至Azure Blob文件拷贝

Airflow DAG实现S3到Azure Blob的每日定时文件拷贝(基于azcopy)

前提条件

  • Airflow节点已安装azcopy工具(推荐v10+版本)
  • Airflow具备访问目标S3存储桶的权限(可通过AWS IAM角色、环境变量AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY或Airflow AWS连接配置)
  • 已确认你提供的Azure Blob存储凭证有效

完整DAG代码

from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta

# 基础参数配置
DEFAULT_ARGS = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

# 替换为你的实际配置
S3_TARGET_PATH = "s3://your-s3-bucket/path/to/target-files/*"  # 可调整匹配规则,比如去掉*同步整个目录
AZURE_ACCOUNT_NAME = "xxx"
AZURE_CONTAINER_NAME = "xxx"
AZURE_SAS_TOKEN = "xxx"  # 与存储账号密钥二选一即可

with DAG(
    's3_to_azure_blob_daily_sync',
    default_args=DEFAULT_ARGS,
    description='Daily sync files from S3 to Azure Blob Storage via azcopy',
    schedule_interval='0 3 * * *',  # 每日凌晨3点触发,可按需求修改Cron表达式
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=['s3', 'azure', 'azcopy', 'data-sync'],
) as dag:

    # 定义azcopy同步任务
    sync_s3_to_azure = BashOperator(
        task_id='sync_s3_to_azure_blob',
        bash_command=f"""
            azcopy sync "{S3_TARGET_PATH}" "https://{AZURE_ACCOUNT_NAME}.blob.core.windows.net/{AZURE_CONTAINER_NAME}?{AZURE_SAS_TOKEN}" \
            --recursive=true \
            --overwrite=ifSourceNewer \
            --log-level=INFO
        """,
        env={
            # 若Airflow未通过IAM角色关联AWS权限,需手动传入密钥(可通过Airflow连接变量引用)
            # 'AWS_ACCESS_KEY_ID': '{{ conn.aws_s3_sync.login }}',
            # 'AWS_SECRET_ACCESS_KEY': '{{ conn.aws_s3_sync.password }}',
        }
    )

    sync_s3_to_azure

关键配置说明

  • 调度时间:schedule_interval='0 3 * * *'对应每日凌晨3点,可参考Cron表达式调整为你需要的特定时间
  • azcopy命令参数:
    • sync:增量同步模式,仅拷贝源端新增/更新的文件,比全量copy更高效
    • --recursive=true:递归处理子目录
    • --overwrite=ifSourceNewer:仅当源文件更新时间晚于目标文件时覆盖,避免重复拷贝
    • 若需全量拷贝,可将sync替换为copy,并调整--overwrite为always
  • Azure认证方式:
    • 使用SAS Token:代码中已采用此方式,直接拼接在Blob容器URL后即可
    • 使用存储账号密钥:可替换命令中的SAS Token部分,改为--account-key={AZURE_STORAGE_ACCOUNT_KEY},同时URL简化为https://{AZURE_ACCOUNT_NAME}.blob.core.windows.net/{AZURE_CONTAINER_NAME}

注意事项

  • 确保azcopy工具在Airflow worker节点的PATH环境变量中,否则需指定完整路径(如/usr/local/bin/azcopy)
  • 若S3路径包含特殊字符,必须用双引号包裹
  • 可添加--log-path=/path/to/azcopy-sync.log参数指定日志路径,方便排查问题
  • 处理大文件时,可调整--concurrency-level=100参数优化同步性能

内容的提问来源于stack exchange,提问作者OneManFunction

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 07:25:26