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

如何用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触发模式。

注意事项

  1. 云函数服务账号需具备GCS读写权限(方案二),并确保SFTP服务器允许云函数的出口IP访问(可通过VPC连接器配置)。
  2. 私钥文件建议通过Secret Manager管理,避免硬编码在代码中。
  3. 超大文件优先选择方案二,利用GCS缓冲避免内存超限;小文件优先方案一,减少中间环节提升效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 01:10:56