如何用Python+GCP无服务器云函数实现SFTP服务器间直接文件传输?
解决方案:基于GCP云函数+asyncssh实现跨SFTP服务器的目录结构保留式文件传输
方案一:源SFTP直接中转到目标SFTP(内存缓冲,适合小文件)
无需经过GCS,直接在云函数内存中完成文件读写,效率更高,适合文件体积不大的场景。
完整代码
import asyncio import asyncssh import stat from flask import jsonify, make_response import functions_framework async def transfer_sftp_to_sftp(source_config, target_config, remote_path): transferred_files = [] # 同时维护源和目标SFTP连接 async with asyncssh.connect( source_config["server_url"], username=source_config["username"], client_keys=[source_config["private_key_path"]], known_hosts=None ) as source_conn, asyncssh.connect( target_config["server_url"], username=target_config["username"], client_keys=[target_config["private_key_path"]], known_hosts=None ) as target_conn: async with source_conn.start_sftp_client() as source_sftp, target_conn.start_sftp_client() as target_sftp: await recursive_transfer(source_sftp, target_sftp, remote_path, remote_path, transferred_files) return transferred_files async def recursive_transfer(source_sftp, target_sftp, source_remote_path, target_base_path, transferred_files): for filename in await source_sftp.listdir(source_remote_path): source_full_path = f"{source_remote_path}/{filename}" target_full_path = f"{target_base_path}/{filename}" attrs = await source_sftp.stat(source_full_path) if stat.S_ISDIR(attrs.permissions): # 在目标端创建对应目录,忽略已存在的情况 try: await target_sftp.stat(target_full_path) except FileNotFoundError: await target_sftp.mkdir(target_full_path) # 递归处理子目录 await recursive_transfer(source_sftp, target_sftp, source_full_path, target_full_path, transferred_files) elif stat.S_ISREG(attrs.permissions): # 分块传输文件,避免内存溢出 async with source_sftp.open(source_full_path, "rb") as source_file: async with target_sftp.open(target_full_path, "wb") as target_file: while chunk := await source_file.read(1024 * 1024): # 按1MB分块读取 await target_file.write(chunk) transferred_files.append(target_full_path) print(f"已传输 {source_full_path} 到 {target_full_path}") @functions_framework.http def main(request): try: # 源SFTP配置 source_config = { "server_url": "your_source_sftp_server", "username": "source_username", "private_key_path": "/path/to/source_private_key" } # 目标SFTP配置 target_config = { "server_url": "your_target_sftp_server", "username": "target_username", "private_key_path": "/path/to/target_private_key" } remote_path = "your_source_remote_directory" # 源SFTP要传输的根目录 transferred_files = asyncio.run(transfer_sftp_to_sftp(source_config, target_config, remote_path)) return make_response( jsonify({ "message": f"文件传输成功,共传输 {len(transferred_files)} 个文件", "transferred_files": transferred_files }), 200 ) except Exception as e: return make_response(jsonify({"error": str(e)}), 500)
核心逻辑说明
- 递归目录同步:遍历源目录时自动判断类型,目录则在目标端创建对应路径,确保层级结构完全一致。
- 分块传输优化:采用1MB分块读写,避免大文件占用过多内存导致云函数崩溃。
- 连接自动管理:用asyncssh上下文管理器维护连接,确保任务完成后自动释放资源。
方案二:基于GCS中转的传输(适合大文件)
针对大文件场景,用GCS作为缓冲,避免云函数内存超限,同时沿用你已有的SFTP→GCS下载逻辑。
完整代码(整合GCS中转)
import asyncio import asyncssh import stat from flask import jsonify, make_response import functions_framework from google.cloud import storage async def download_from_sftp_to_gcs(server_config, remote_path, bucket_name): storage_client = storage.Client() bucket = storage_client.bucket(bucket_name) downloaded_files = [] async with asyncssh.connect( server_config["server_url"], username=server_config["username"], client_keys=[server_config["private_key_path"]], known_hosts=None ) as conn: async with conn.start_sftp_client() as sftp: await recursive_download(sftp, remote_path, bucket, downloaded_files) return downloaded_files async def recursive_download(sftp, remote_path, bucket, downloaded_files): for filename in await sftp.listdir(remote_path): remote_file = f"{remote_path}/{filename}" attrs = await sftp.stat(remote_file) if stat.S_ISREG(attrs.permissions): async with sftp.open(remote_file) as remote_file_obj: file_data = await remote_file_obj.read() blob = bucket.blob(remote_file) blob.upload_from_string(file_data) downloaded_files.append(remote_file) print(f"已下载 {remote_file} 到GCS") elif stat.S_ISDIR(attrs.permissions): await recursive_download(sftp, remote_file, bucket, downloaded_files) async def upload_from_gcs_to_sftp(server_config, bucket_name, remote_path): storage_client = storage.Client() bucket = storage_client.bucket(bucket_name) uploaded_files = [] async with asyncssh.connect( server_config["server_url"], username=server_config["username"], client_keys=[server_config["private_key_path"]], known_hosts=None ) as conn: async with conn.start_sftp_client() as sftp: # 获取GCS中对应路径的所有文件 blobs = bucket.list_blobs(prefix=remote_path) for blob in blobs: if blob.name.endswith("/"): continue # 跳过GCS的虚拟目录占位 # 递归创建目标SFTP的父目录 target_dir = "/".join(blob.name.split("/")[:-1]) try: await sftp.stat(target_dir) except FileNotFoundError: await sftp.mkdir(target_dir, parents=True) # 读取GCS文件并写入目标SFTP file_data = blob.download_as_bytes() async with sftp.open(blob.name, "wb") as target_file: await target_file.write(file_data) uploaded_files.append(blob.name) print(f"已从GCS上传 {blob.name} 到SFTP") return uploaded_files @functions_framework.http def main(request): try: source_config = { "server_url": "your_source_sftp_server", "username": "source_username", "private_key_path": "/path/to/source_private_key" } target_config = { "server_url": "your_target_sftp_server", "username": "target_username", "private_key_path": "/path/to/target_private_key" } remote_path = "your_source_remote_directory" bucket_name = "your_gcs_bucket_name" # 第一步:源SFTP→GCS downloaded_files = asyncio.run(download_from_sftp_to_gcs(source_config, remote_path, bucket_name)) # 第二步:GCS→目标SFTP uploaded_files = asyncio.run(upload_from_gcs_to_sftp(target_config, bucket_name, remote_path)) return make_response( jsonify({ "message": f"中转传输完成,共下载 {len(downloaded_files)} 个文件,上传 {len(uploaded_files)} 个文件", "downloaded_files": downloaded_files, "uploaded_files": uploaded_files }), 200 ) except Exception as e: return make_response(jsonify({"error": str(e)}), 500)
关键优化点
- GCS路径保留:通过Blob的
name属性完整保留原SFTP的目录结构,上传时自动同步到目标SFTP。 - 异步API适配:所有SFTP操作使用asyncssh异步接口,适配GCP云函数的HTTP触发模式。
注意事项
- 云函数服务账号需具备GCS读写权限(方案二),并确保SFTP服务器允许云函数的出口IP访问(可通过VPC连接器配置)。
- 私钥文件建议通过Secret Manager管理,避免硬编码在代码中。
- 超大文件优先选择方案二,利用GCS缓冲避免内存超限;小文件优先方案一,减少中间环节提升效率。
内容的提问来源于stack exchange,提问作者Safiul Alam
相关产品推荐
相关产品推荐

