请求协助:在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
相关产品推荐
相关产品推荐

