如何使用Airflow实现GCS到远程SFTP服务器的文件传输
GCS到远程SFTP无本地中转的Airflow实现方案
完全可以实现,全程不需要部署常驻Unix中转服务器,以下是3种生产可用的落地路径,附实现参考:
方案1:Airflow原生Provider流传输(最轻量,无额外资源依赖)
这个方案不需要拉起额外云主机或者容器实例,直接通过Airflow官方Hook的流能力在内存中完成数据传输,文件不会全量落盘到Worker本地存储。
- 前置依赖:在Airflow环境安装两个官方Provider包:
apache-airflow-providers-google、apache-airflow-providers-sftp - 核心逻辑:跳过“下载GCS文件到本地磁盘”的步骤,直接获取GCS文件的可读流对象,通过SFTP客户端的上传接口把流直接写入远程SFTP服务器,全程数据仅在内存中分片流转。
- 参考实现代码:
from airflow import DAG from airflow.decorators import task from airflow.providers.google.cloud.hooks.gcs import GCSHook from airflow.providers.sftp.hooks.sftp import SFTPHook from datetime import datetime with DAG( dag_id="gcs_to_sftp_direct_stream", start_date=datetime(2024, 1, 1), schedule=None, catchup=False ) as dag: @task def run_transfer(): # 替换为自己环境配置的连接ID gcs_hook = GCSHook(gcp_conn_id="gcs_prod_conn") sftp_hook = SFTPHook(ssh_conn_id="sftp_remote_conn") # 替换为实际的源和目标路径 source_bucket = "core-business-data-bucket" source_gcs_path = "export/2024/transaction_data.parquet" target_sftp_path = "/data/upload/transaction_data_2024.parquet" # 获取GCS文件流,不写入本地磁盘 with gcs_hook.provide_file( bucket_name=source_bucket, object_name=source_gcs_path ) as gcs_stream: # 直接将流写入远程SFTP sftp_hook.store_file( remote_full_path=target_sftp_path, local_full_path_or_buffer=gcs_stream ) run_transfer()
适配场景:单文件大小在10G以内的常规传输,资源占用极低,配置完成后几乎不需要额外运维。如果传输10G以上大文件,给Worker预留2G以上空闲内存即可,避免流传输缓冲占用过高。
方案2:GCP Cloud Run Job托管传输(不占用Airflow Worker资源)
如果不想让传输流量、计算占用Airflow Worker资源,可以用GCP托管的Cloud Run Job承载传输逻辑,Airflow只负责任务触发和状态校验,全程没有自维护的中转节点。
- 实现步骤:
- 把和方案1一致的流传输逻辑打包成轻量Docker镜像,镜像仅需安装
google-cloud-storage、paramiko两个依赖,全程不做本地落盘操作 - 给Cloud Run Job配置拥有GCS对象读取权限的服务账号,SFTP连接密钥存在GCP Secret Manager,Job启动时自动拉取
- Airflow侧通过
CloudRunExecuteJobOperator触发Job运行,等待Job执行完成后标记任务成功,传输过程中所有计算、带宽资源都由Cloud Run承载,不占用Airflow集群资源
- 把和方案1一致的流传输逻辑打包成轻量Docker镜像,镜像仅需安装
- 优势:Cloud Run是Serverless服务,任务跑完自动释放资源,没有常驻成本,单文件最大支持到100G级别传输,稳定性更高。
方案3:临时K8s Pod隔离传输(适配K8s部署的Airflow集群)
如果Airflow部署在Kubernetes集群上,可以直接用KubernetesPodOperator拉起临时Pod运行传输逻辑,任务跑完Pod自动销毁,不需要部署常驻中转服务器。
- 实现要点:
- 用轻量Python基础镜像构建传输镜像,内置流传输逻辑,不需要额外组件
- 可以给传输Pod单独配置CPU、内存、带宽限制,和Airflow核心组件的资源完全隔离,不会出现大文件传输堵死Worker的问题
- SFTP密钥、GCP凭证可以通过K8s Secret挂载到临时Pod,不需要硬编码在镜像里
避坑提醒:不要使用
GCSToLocalFilesystemOperator先把GCS文件下载到本地路径再调用SFTP上传,这种方式会把文件全量落盘到执行节点本地存储,不符合无本地中转的要求。
内容的提问来源于stack exchange,提问作者Khilesh Chauhan
相关产品推荐
相关产品推荐

