如何通过Airflow DAG将S3文件复制到EC2挂载的EBS指定目录
实现S3到EC2挂载EBS的文件复制Airflow DAG
以下是两种实用的实现方案,可根据你的环境选择:
方法1:使用Airflow原生Hook(推荐)
通过Airflow的S3Hook将文件下载到worker节点,再通过SSHHook上传到EC2的EBS挂载目录,无需在EC2预安装awscli。
前提条件
- Airflow已配置具备S3读权限的AWS连接(默认
aws_default) - Airflow已配置可连接目标EC2的SSH连接(命名为
ec2_ssh_conn),且SSH用户对/usr/local/myfiles有写入权限 - EC2安全组允许Airflow worker的IP通过22端口访问
DAG代码示例
from airflow import DAG from airflow.providers.amazon.aws.hooks.s3 import S3Hook from airflow.providers.ssh.hooks.ssh import SSHHook from airflow.operators.python import PythonOperator from datetime import datetime, timedelta import os def copy_s3_to_ec2_ebs(): # 初始化Hook s3_hook = S3Hook(aws_conn_id='aws_default') ssh_hook = SSHHook(ssh_conn_id='ec2_ssh_conn') # 配置路径参数 s3_bucket = 'your-s3-bucket-name' s3_object_key = 'target/file/path/on/s3.txt' local_temp_file = f'/tmp/{os.path.basename(s3_object_key)}' ebs_target_path = f'/usr/local/myfiles/{os.path.basename(s3_object_key)}' # 从S3下载到worker临时目录 s3_hook.download_file( bucket_name=s3_bucket, key=s3_object_key, local_path=local_temp_file ) # 通过SFTP上传到EC2的EBS目录 with ssh_hook.get_conn() as ssh_conn: sftp_client = ssh_conn.open_sftp() try: sftp_client.put(local_temp_file, ebs_target_path) finally: sftp_client.close() # 清理worker临时文件 os.remove(local_temp_file) default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'retries': 1, 'retry_delay': timedelta(minutes=5), } with DAG( 's3_to_ec2_ebs_sync', default_args=default_args, description='Sync files from S3 to EC2 mounted EBS volume', schedule_interval='@daily', catchup=False, ) as dag: copy_task = PythonOperator( task_id='transfer_s3_to_ec2', python_callable=copy_s3_to_ec2_ebs, ) copy_task
方法2:使用SSHOperator配合awscli(适合EC2直连S3场景)
如果EC2实例通过IAM角色拥有S3访问权限,可直接在EC2上执行aws s3 cp命令,通过Airflow远程触发。
前提条件
- EC2实例的IAM角色具备S3读权限
- EC2已安装并配置awscli(通过IAM角色自动认证)
- Airflow已配置EC2的SSH连接
DAG代码示例
from airflow import DAG from airflow.providers.ssh.operators.ssh import SSHOperator from datetime import datetime, timedelta default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'retries': 1, 'retry_delay': timedelta(minutes=5), } with DAG( 's3_to_ec2_ebs_cli_sync', default_args=default_args, description='Copy files from S3 to EC2 EBS via awscli', schedule_interval='@daily', catchup=False, ) as dag: single_file_copy = SSHOperator( task_id='copy_single_file', ssh_conn_id='ec2_ssh_conn', command='aws s3 cp s3://your-bucket/path/file.txt /usr/local/myfiles/' ) # 如需同步整个目录,替换为sync命令 dir_sync = SSHOperator( task_id='sync_directory', ssh_conn_id='ec2_ssh_conn', command='aws s3 sync s3://your-bucket/source-dir/ /usr/local/myfiles/' ) single_file_copy >> dir_sync
关键注意事项
- 权限校验:确保各环节的权限配置正确(S3读权限、EC2目录写入权限、SSH访问权限)
- 大文件处理:大文件推荐用方法2,避免占用Airflow worker磁盘空间
- 调度调整:根据实际需求修改
schedule_interval参数
内容的提问来源于stack exchange,提问作者Ab_sin
相关产品推荐
相关产品推荐

